From f8fc9524e5d1f6f0f467d542ef532afa1d5de560 Mon Sep 17 00:00:00 2001 From: y-oga Date: Fri, 5 Jun 2026 16:35:20 +0900 Subject: [PATCH 1/3] =?UTF-8?q?fix:=20DynamoDB=E3=81=AE=E3=82=A4=E3=83=99?= =?UTF-8?q?=E3=83=B3=E3=83=88=E5=8F=96=E5=BE=97=E9=A0=86=E3=82=92payload.t?= =?UTF-8?q?imestamp=E3=81=AB=E7=B5=B1=E4=B8=80?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit DynamoDB backend の FindBySessionID は created_at(サーバ受信時刻)で整列 していたが、sqlite/postgres/turso/memory は transcript 行の payload.timestamp(イベント時刻)で整列しており食い違っていた。同期・順序 どおりの配送では両者がほぼ一致するため顕在化しなかったが、遅延や順不同の 取り込みでは DynamoDB だけ表示順が崩れる。 sort_key が重複排除のため UUID ベースになった際に時系列順が失われ、その 代替として created_at ソートが入った経緯による(コード内コメントにも記載)。 他バックエンドに揃える: - event.go: sort.SliceStable + payload.timestamp 整列に変更。欠落時は created_at にフォールバック(RFC3339Nano → "...000Z" → CreatedAt)で、 sqlite/turso の getTimestampFromPayload と同一ロジック。 - testsuite: created_at の順を payload.timestamp の逆順に仕込んだ共有順序 テストを追加し、全バックエンドの整列契約を固定。修正前は DynamoDB での み fail、他バックエンドでは pass。 Co-Authored-By: Claude Opus 4.8 --- server/internal/repository/dynamodb/event.go | 27 ++++++++-- .../repository/testsuite/event_suite.go | 54 +++++++++++++++++++ 2 files changed, 78 insertions(+), 3 deletions(-) diff --git a/server/internal/repository/dynamodb/event.go b/server/internal/repository/dynamodb/event.go index c53e8ea..21f62a2 100644 --- a/server/internal/repository/dynamodb/event.go +++ b/server/internal/repository/dynamodb/event.go @@ -135,14 +135,35 @@ func (r *EventRepository) FindBySessionID(ctx context.Context, sessionID string) events[i] = r.itemToEvent(&item) } - // Sort by created_at since sort_key is now uuid-based and doesn't preserve chronological order - sort.Slice(events, func(i, j int) bool { - return events[i].CreatedAt.Before(events[j].CreatedAt) + // The sort_key is uuid-based and does not preserve chronological order, so the + // ordering has to be reconstructed here. Sort by the transcript line's own + // payload.timestamp (falling back to created_at) to match the sqlite/postgres/ + // turso/memory backends; created_at is the server receive time and diverges from + // the event time under out-of-order delivery. + sort.SliceStable(events, func(i, j int) bool { + return eventTimestamp(events[i]).Before(eventTimestamp(events[j])) }) return events, nil } +// eventTimestamp returns the ordering key for an event: the transcript line's own +// payload.timestamp when present and parseable, otherwise the server-side CreatedAt. +// This mirrors getTimestampFromPayload in the sqlite/turso backends so that every +// backend orders a session's events identically. +func eventTimestamp(e *domain.Event) time.Time { + if ts, ok := e.Payload["timestamp"].(string); ok { + if parsed, err := time.Parse(time.RFC3339Nano, ts); err == nil { + return parsed + } + // Fallback for timestamps without timezone offset. + if parsed, err := time.Parse("2006-01-02T15:04:05.000Z", ts); err == nil { + return parsed + } + } + return e.CreatedAt +} + func (r *EventRepository) CountBySessionID(ctx context.Context, sessionID string) (int, error) { keyCond := expression.Key("session_id").Equal(expression.Value(sessionID)) diff --git a/server/internal/repository/testsuite/event_suite.go b/server/internal/repository/testsuite/event_suite.go index 6a4851f..a455b78 100644 --- a/server/internal/repository/testsuite/event_suite.go +++ b/server/internal/repository/testsuite/event_suite.go @@ -228,6 +228,60 @@ func (s *EventRepositorySuite) TestFindBySessionID_ChronologicalOrder() { s.WithinDuration(baseTime.Add(500*time.Millisecond), events[4].CreatedAt, time.Microsecond, "Last event should have latest timestamp") } +// TestFindBySessionID_OrdersByPayloadTimestamp verifies that events are ordered by +// the transcript line's own payload.timestamp, NOT by created_at (server receive time). +// +// This guards against backend drift: SQLite/Postgres/turso/memory already sort by +// payload.timestamp, while DynamoDB used to sort by created_at. Under asynchronous / +// out-of-order delivery the two disagree, so we deliberately insert events whose +// created_at order is the REVERSE of their payload.timestamp order. +func (s *EventRepositorySuite) TestFindBySessionID_OrdersByPayloadTimestamp() { + ctx := context.Background() + + sessionID := s.createTestSession("event-payload-ts-order") + if sessionID == "" { + s.T().Skip("SessionRepo not available, skipping test") + } + + // Second precision + UTC "Z" keeps the timestamp string lexically ordered, + // which Postgres relies on (ORDER BY payload->>'timestamp'), while Go-based + // backends parse it with RFC3339Nano. Fixed width avoids fractional-trim pitfalls. + baseTime := time.Now().UTC().Truncate(time.Second) + + const n = 5 + expectedTimestamps := make([]string, n) + for k := 0; k < n; k++ { + // payload.timestamp ascends with k; created_at descends with k (reverse order). + ts := baseTime.Add(time.Duration(k) * time.Second).Format(time.RFC3339Nano) + expectedTimestamps[k] = ts + event := &domain.Event{ + SessionID: sessionID, + UUID: uuid.New().String(), + EventType: "message", + Payload: map[string]interface{}{ + "timestamp": ts, + "marker": k, + }, + CreatedAt: baseTime.Add(time.Duration(n-k) * time.Hour), + } + err := s.Repo.Create(ctx, event) + s.Require().NoError(err) + } + + events, err := s.Repo.FindBySessionID(ctx, sessionID) + s.Require().NoError(err) + s.Require().Len(events, n) + + // Returned order must follow payload.timestamp ascending, i.e. the insertion (created_at) + // order reversed. If a backend sorted by created_at, this would come back descending and fail. + for i := 0; i < n; i++ { + ts, ok := events[i].Payload["timestamp"].(string) + s.Require().True(ok, "event[%d] should carry a string payload.timestamp", i) + s.Equal(expectedTimestamps[i], ts, + "events should be ordered by payload.timestamp ascending (position %d)", i) + } +} + func (s *EventRepositorySuite) TestFindBySessionID_Empty() { ctx := context.Background() From da3c0acad0bc73da810adefc9a9e8eca822e0b2c Mon Sep 17 00:00:00 2001 From: y-oga Date: Fri, 5 Jun 2026 16:45:57 +0900 Subject: [PATCH 2/3] =?UTF-8?q?style:=20=E3=82=B3=E3=83=A1=E3=83=B3?= =?UTF-8?q?=E3=83=88=E3=82=92=E5=91=A8=E8=BE=BA=E3=82=B3=E3=83=BC=E3=83=89?= =?UTF-8?q?=E3=81=AE=E5=AF=86=E5=BA=A6=E3=81=AB=E5=90=88=E3=82=8F=E3=81=9B?= =?UTF-8?q?=E3=81=A6=E7=B0=A1=E6=BD=94=E5=8C=96?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 追加した sort/helper/test のコメントが周辺(sqlite の getTimestampFromPayload は doc 無し、既存テストもメソッド doc 無し)より過剰だったため house style に 合わせて簡潔化。ロジックは不変、テスト green。 Co-Authored-By: Claude Opus 4.8 --- server/internal/repository/dynamodb/event.go | 13 +++---------- server/internal/repository/testsuite/event_suite.go | 9 ++------- 2 files changed, 5 insertions(+), 17 deletions(-) diff --git a/server/internal/repository/dynamodb/event.go b/server/internal/repository/dynamodb/event.go index 21f62a2..350ae39 100644 --- a/server/internal/repository/dynamodb/event.go +++ b/server/internal/repository/dynamodb/event.go @@ -135,11 +135,8 @@ func (r *EventRepository) FindBySessionID(ctx context.Context, sessionID string) events[i] = r.itemToEvent(&item) } - // The sort_key is uuid-based and does not preserve chronological order, so the - // ordering has to be reconstructed here. Sort by the transcript line's own - // payload.timestamp (falling back to created_at) to match the sqlite/postgres/ - // turso/memory backends; created_at is the server receive time and diverges from - // the event time under out-of-order delivery. + // sort_key is uuid-based, so chronological order must be reconstructed here. + // Match the other backends: sort by payload.timestamp, not created_at. sort.SliceStable(events, func(i, j int) bool { return eventTimestamp(events[i]).Before(eventTimestamp(events[j])) }) @@ -147,16 +144,12 @@ func (r *EventRepository) FindBySessionID(ctx context.Context, sessionID string) return events, nil } -// eventTimestamp returns the ordering key for an event: the transcript line's own -// payload.timestamp when present and parseable, otherwise the server-side CreatedAt. -// This mirrors getTimestampFromPayload in the sqlite/turso backends so that every -// backend orders a session's events identically. func eventTimestamp(e *domain.Event) time.Time { if ts, ok := e.Payload["timestamp"].(string); ok { if parsed, err := time.Parse(time.RFC3339Nano, ts); err == nil { return parsed } - // Fallback for timestamps without timezone offset. + // Try parsing without timezone if parsed, err := time.Parse("2006-01-02T15:04:05.000Z", ts); err == nil { return parsed } diff --git a/server/internal/repository/testsuite/event_suite.go b/server/internal/repository/testsuite/event_suite.go index a455b78..e51579c 100644 --- a/server/internal/repository/testsuite/event_suite.go +++ b/server/internal/repository/testsuite/event_suite.go @@ -228,13 +228,8 @@ func (s *EventRepositorySuite) TestFindBySessionID_ChronologicalOrder() { s.WithinDuration(baseTime.Add(500*time.Millisecond), events[4].CreatedAt, time.Microsecond, "Last event should have latest timestamp") } -// TestFindBySessionID_OrdersByPayloadTimestamp verifies that events are ordered by -// the transcript line's own payload.timestamp, NOT by created_at (server receive time). -// -// This guards against backend drift: SQLite/Postgres/turso/memory already sort by -// payload.timestamp, while DynamoDB used to sort by created_at. Under asynchronous / -// out-of-order delivery the two disagree, so we deliberately insert events whose -// created_at order is the REVERSE of their payload.timestamp order. +// TestFindBySessionID_OrdersByPayloadTimestamp ensures events are returned ordered by +// payload.timestamp (event time), not created_at — guarding the cross-backend contract. func (s *EventRepositorySuite) TestFindBySessionID_OrdersByPayloadTimestamp() { ctx := context.Background() From bc93b1c689da7fd3e3abda2581f95e8c6b19c53e Mon Sep 17 00:00:00 2001 From: y-oga Date: Fri, 5 Jun 2026 17:07:57 +0900 Subject: [PATCH 3/3] =?UTF-8?q?refactor:=20=E3=83=98=E3=83=AB=E3=83=91?= =?UTF-8?q?=E5=90=8D=E3=82=92=E4=BB=96=E3=83=90=E3=83=83=E3=82=AF=E3=82=A8?= =?UTF-8?q?=E3=83=B3=E3=83=89=E3=81=AB=E5=90=88=E3=82=8F=E3=81=9B=E3=81=A6?= =?UTF-8?q?getTimestampFromPayload=E3=81=AB=E7=B5=B1=E4=B8=80?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit DynamoDB に追加した整列ヘルパを eventTimestamp から、memory/sqlite/turso が 使う getTimestampFromPayload に rename。機能は不変、命名を揃えるだけ。 Co-Authored-By: Claude Opus 4.8 --- server/internal/repository/dynamodb/event.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/server/internal/repository/dynamodb/event.go b/server/internal/repository/dynamodb/event.go index 350ae39..5801e30 100644 --- a/server/internal/repository/dynamodb/event.go +++ b/server/internal/repository/dynamodb/event.go @@ -138,13 +138,13 @@ func (r *EventRepository) FindBySessionID(ctx context.Context, sessionID string) // sort_key is uuid-based, so chronological order must be reconstructed here. // Match the other backends: sort by payload.timestamp, not created_at. sort.SliceStable(events, func(i, j int) bool { - return eventTimestamp(events[i]).Before(eventTimestamp(events[j])) + return getTimestampFromPayload(events[i]).Before(getTimestampFromPayload(events[j])) }) return events, nil } -func eventTimestamp(e *domain.Event) time.Time { +func getTimestampFromPayload(e *domain.Event) time.Time { if ts, ok := e.Payload["timestamp"].(string); ok { if parsed, err := time.Parse(time.RFC3339Nano, ts); err == nil { return parsed