Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 10 additions & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -21,3 +21,13 @@ Backward compatibility is not a primary constraint in the current development st

1. Make small, focused changes.
2. Run format/lint/tests relevant to the change.

## Commit Format

Use Conventional Commit format for commit subjects and PR titles:

`type(scope): short summary`

Allowed types: feat, fix, perf, docs, refactor, test, build, ci, chore, revert

Squash merges should preserve the PR title as the resulting commit subject.
10 changes: 10 additions & 0 deletions CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,16 @@ make swagger # Regenerate Swagger docs

**Build tags:** E2E tests require `-tags=e2e`, integration tests require `-tags=integration`, contract tests require `-tags=contract`. The Makefile handles this automatically.

## Commit And PR Title Format

Use Conventional Commit format for commit subjects and PR titles:

`type(scope): short summary`

Allowed types: feat, fix, perf, docs, refactor, test, build, ci, chore, revert

Prefer squash-and-merge to keep the merged commit subject aligned with the PR title.

## Error Handling

- All errors returned to clients must be instances of `core.GatewayError`.
Expand Down
42 changes: 18 additions & 24 deletions internal/auditlog/auditlog_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import (
"encoding/json"
"fmt"
"gomodel/internal/core"
"gomodel/internal/streaming"
"io"
"net/http"
"net/http/httptest"
Expand Down Expand Up @@ -705,7 +706,7 @@ func TestIsModelInteractionPath(t *testing.T) {
}
}

func TestStreamLogWrapper(t *testing.T) {
func TestStreamLogObserver(t *testing.T) {
// Create a mock stream with content
streamContent := `data: {"id":"chatcmpl-123","choices":[{"delta":{"content":"Hello"}}]}

Expand All @@ -724,7 +725,6 @@ data: [DONE]
FlushInterval: 100 * time.Millisecond,
}
logger := NewLogger(store, cfg)
defer logger.Close()

entry := &LogEntry{
ID: "test-entry",
Expand All @@ -733,44 +733,38 @@ data: [DONE]
Data: &LogData{},
}

// Wrap the stream
wrapper := NewStreamLogWrapper(stream, logger, entry, "/v1/chat/completions")
observedStream := streaming.NewObservedSSEStream(
stream,
NewStreamLogObserver(logger, entry, "/v1/chat/completions"),
)

// Read all content
var buf bytes.Buffer
_, err := io.Copy(&buf, wrapper)
_, err := io.Copy(&buf, observedStream)
if err != nil {
t.Fatalf("failed to read stream: %v", err)
}

// Close wrapper to trigger logging
if err := wrapper.Close(); err != nil {
t.Fatalf("failed to close wrapper: %v", err)
// Close stream to trigger logging
if err := observedStream.Close(); err != nil {
t.Fatalf("failed to close stream: %v", err)
}
if err := logger.Close(); err != nil {
t.Fatalf("failed to close logger: %v", err)
}

// Wait for async write
time.Sleep(200 * time.Millisecond)

// Verify entry was logged
if len(store.getEntries()) != 1 {
t.Errorf("expected 1 entry, got %d", len(store.getEntries()))
}
}

func TestWrapStreamForLogging(t *testing.T) {
stream := io.NopCloser(strings.NewReader("test"))

// Test with nil logger
result := WrapStreamForLogging(stream, nil, nil, "/v1/chat/completions")
if result != stream {
t.Error("expected original stream with nil logger")
func TestNewStreamLogObserverNilInputs(t *testing.T) {
if observer := NewStreamLogObserver(nil, &LogEntry{}, "/v1/chat/completions"); observer != nil {
t.Error("expected nil observer with nil logger")
}

// Test with disabled logger
noopLogger := &NoopLogger{}
result = WrapStreamForLogging(stream, noopLogger, &LogEntry{}, "/v1/chat/completions")
if result != stream {
t.Error("expected original stream with disabled logger")
if observer := NewStreamLogObserver(&NoopLogger{}, nil, "/v1/chat/completions"); observer != nil {
t.Error("expected nil observer with nil entry")
}
}

Expand Down
5 changes: 3 additions & 2 deletions internal/auditlog/constants.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ const (
MaxBodyCapture = 1024 * 1024

// MaxContentCapture is the maximum size of accumulated streaming content (1MB).
// Used by StreamLogWrapper to limit reconstructed response body size.
// Used by the stream observer to limit reconstructed response body size.
MaxContentCapture = 1024 * 1024

// BatchFlushThreshold is the number of entries that triggers an immediate flush.
Expand All @@ -27,6 +27,7 @@ const (
LogEntryKey contextKey = "auditlog_entry"

// LogEntryStreamingKey is the context key for marking a request as streaming.
// When true, the middleware skips logging (StreamLogWrapper handles it instead).
// When true, the middleware skips logging because the stream observer path
// handles streaming audit logging.
LogEntryStreamingKey contextKey = "auditlog_entry_streaming"
)
4 changes: 2 additions & 2 deletions internal/auditlog/middleware.go
Original file line number Diff line number Diff line change
Expand Up @@ -149,7 +149,7 @@ func Middleware(logger LoggerInterface) echo.MiddlewareFunc {
}
}

// Write log entry asynchronously (skip if streaming - StreamLogWrapper handles it)
// Write log entry asynchronously (skip if streaming - the stream observer path handles it)
if !IsEntryMarkedAsStreaming(c) {
logger.Write(entry)
}
Expand Down Expand Up @@ -262,7 +262,7 @@ type responseBodyCapture struct {
body *bytes.Buffer
truncated bool
// shouldCapture allows middleware to stop buffering once the request is
// known to be streaming. Streaming responses are handled by StreamLogWrapper.
// known to be streaming. Streaming responses are handled by the stream observer path.
shouldCapture func() bool
}

Expand Down
152 changes: 152 additions & 0 deletions internal/auditlog/stream_observer.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,152 @@
package auditlog

import (
"strings"
"time"
)

// StreamLogObserver reconstructs stream metadata and optional response bodies
// from parsed SSE JSON payloads.
type StreamLogObserver struct {
logger LoggerInterface
entry *LogEntry
builder *streamResponseBuilder
logBodies bool
closed bool
startTime time.Time
}

func NewStreamLogObserver(logger LoggerInterface, entry *LogEntry, path string) *StreamLogObserver {
if logger == nil || entry == nil {
return nil
}

logBodies := logger.Config().LogBodies
var builder *streamResponseBuilder
if logBodies {
builder = &streamResponseBuilder{
IsResponsesAPI: strings.HasPrefix(path, "/v1/responses"),
}
}

return &StreamLogObserver{
logger: logger,
entry: entry,
builder: builder,
logBodies: logBodies,
startTime: entry.Timestamp,
}
}

func (o *StreamLogObserver) OnJSONEvent(event map[string]interface{}) {
if !o.logBodies || o.builder == nil {
return
}
if o.builder.IsResponsesAPI {
o.parseResponsesAPIEvent(event)
return
}
o.parseChatCompletionEvent(event)
}

func (o *StreamLogObserver) OnStreamClose() {
if o.closed {
return
}
o.closed = true

if o.entry != nil && !o.startTime.IsZero() {
o.entry.DurationNs = time.Since(o.startTime).Nanoseconds()
}

if o.logBodies && o.builder != nil && o.entry != nil && o.entry.Data != nil {
if o.builder.IsResponsesAPI {
o.entry.Data.ResponseBody = o.builder.buildResponsesAPIResponse()
} else {
o.entry.Data.ResponseBody = o.builder.buildChatCompletionResponse()
}
o.entry.Data.ResponseBodyTooBigToHandle = o.builder.truncated
}

if o.logger != nil && o.entry != nil {
o.logger.Write(o.entry)
}
}

func (o *StreamLogObserver) parseChatCompletionEvent(event map[string]interface{}) {
if o.builder == nil {
return
}

if o.builder.ID == "" {
if id, ok := event["id"].(string); ok {
o.builder.ID = id
}
if model, ok := event["model"].(string); ok {
o.builder.Model = model
}
if created, ok := event["created"].(float64); ok {
o.builder.Created = int64(created)
Comment on lines +88 to +89

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🧹 Nitpick | 🔵 Trivial

Consider precision loss for large Unix timestamps.

JSON numbers unmarshal as float64, which can represent integers exactly only up to 2^53. While current Unix timestamps (~1.7×10^9) are safe, the cast int64(created) would lose precision for nanosecond timestamps or far-future dates. This is likely fine for the created field (seconds since epoch), but worth noting if this pattern is reused elsewhere.

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@internal/auditlog/stream_observer.go` around lines 88 - 89, The code takes
created from event as a float64 and casts to int64 which can lose precision for
very large timestamps; switch JSON decoding to use json.Decoder.UseNumber (or
ensure the unmarshalling produces json.Number) and here detect json.Number in
event["created"], calling .Int64() (or strconv.ParseInt on its String()) to
safely convert to int64 before assigning to o.builder.Created; alternatively, if
event may contain a string timestamp, handle that case by parsing with
strconv.ParseInt — update the conversion logic around the created variable and
the assignment to o.builder.Created accordingly.

}
}
Comment on lines +81 to +91

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🟡 Minor

Decouple model/created extraction from id initialization.

At Line 81, model and created are only parsed while ID is empty. If ID arrives before those fields, they can remain unset.

💡 Proposed fix
-	if o.builder.ID == "" {
-		if id, ok := event["id"].(string); ok {
-			o.builder.ID = id
-		}
-		if model, ok := event["model"].(string); ok {
-			o.builder.Model = model
-		}
-		if created, ok := event["created"].(float64); ok {
-			o.builder.Created = int64(created)
-		}
-	}
+	if o.builder.ID == "" {
+		if id, ok := event["id"].(string); ok {
+			o.builder.ID = id
+		}
+	}
+	if o.builder.Model == "" {
+		if model, ok := event["model"].(string); ok {
+			o.builder.Model = model
+		}
+	}
+	if o.builder.Created == 0 {
+		if created, ok := event["created"].(float64); ok {
+			o.builder.Created = int64(created)
+		}
+	}
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@internal/auditlog/stream_observer.go` around lines 81 - 91, The code
currently only assigns o.builder.Model and o.builder.Created when o.builder.ID
== "" which can leave Model/Created unset if ID was present first; change the
logic so extracting "model" and "created" from the event map is done
independently of the ID check: keep the existing ID assignment guarded by if
o.builder.ID == "" and the id type assertion, but move the model and created
parsing outside that block (or into their own if checks) so that if model
(string) or created (float64) exist in event they are assigned to
o.builder.Model and o.builder.Created regardless of o.builder.ID.


if choices, ok := event["choices"].([]interface{}); ok && len(choices) > 0 {
if choice, ok := choices[0].(map[string]interface{}); ok {
if fr, ok := choice["finish_reason"].(string); ok && fr != "" {
o.builder.FinishReason = fr
}

if delta, ok := choice["delta"].(map[string]interface{}); ok {
if role, ok := delta["role"].(string); ok {
o.builder.Role = role
}
if content, ok := delta["content"].(string); ok && content != "" {
o.appendContent(content)
}
}
}
}
}

func (o *StreamLogObserver) parseResponsesAPIEvent(event map[string]interface{}) {
if o.builder == nil {
return
}

eventType, _ := event["type"].(string)
switch eventType {
case "response.created", "response.completed", "response.done":
if resp, ok := event["response"].(map[string]interface{}); ok {
if id, ok := resp["id"].(string); ok {
o.builder.ResponseID = id
}
if status, ok := resp["status"].(string); ok {
o.builder.Status = status
}
if model, ok := resp["model"].(string); ok {
o.builder.Model = model
}
if createdAt, ok := resp["created_at"].(float64); ok {
o.builder.CreatedAt = int64(createdAt)
}
}
case "response.output_text.delta":
if delta, ok := event["delta"].(string); ok && delta != "" {
o.appendContent(delta)
}
}
}

func (o *StreamLogObserver) appendContent(content string) {
if o.builder == nil || o.builder.truncated || o.builder.contentLen >= MaxContentCapture {
return
}
Comment on lines +141 to +143

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🟠 Major

Set truncation flag when content is dropped at capacity.

At Line 141, extra chunks are dropped once capacity is reached, but truncated may remain false. That can make ResponseBodyTooBigToHandle incorrect on close.

💡 Proposed fix
 func (o *StreamLogObserver) appendContent(content string) {
-	if o.builder == nil || o.builder.truncated || o.builder.contentLen >= MaxContentCapture {
+	if o.builder == nil || o.builder.truncated || content == "" {
+		return
+	}
+	if o.builder.contentLen >= MaxContentCapture {
+		o.builder.truncated = true
 		return
 	}
 
 	remaining := MaxContentCapture - o.builder.contentLen
 	if len(content) > remaining {
 		content = content[:remaining]
 		o.builder.truncated = true
 	}
 	o.builder.Content.WriteString(content)
 	o.builder.contentLen += len(content)
 }

Also applies to: 145-152

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@internal/auditlog/stream_observer.go` around lines 141 - 143, The code drops
extra chunks once o.builder.contentLen >= MaxContentCapture but doesn't set
o.builder.truncated, causing ResponseBodyTooBigToHandle to be incorrect; update
the early-return branches that check "if o.builder == nil || o.builder.truncated
|| o.builder.contentLen >= MaxContentCapture" (and the similar branch in the
145-152 region) to set o.builder.truncated = true before returning when
contentLen >= MaxContentCapture (and ensure you check builder != nil first), so
the truncated flag accurately reflects that content was dropped.


remaining := MaxContentCapture - o.builder.contentLen
if len(content) > remaining {
content = content[:remaining]
o.builder.truncated = true
}
o.builder.Content.WriteString(content)
o.builder.contentLen += len(content)
}
Loading
Loading