Skip to content

Commit 7c975f1

Browse files
fix: filter duckdb usage window before dedup and use shared trend counting
1 parent f7f7a82 commit 7c975f1

3 files changed

Lines changed: 123 additions & 54 deletions

File tree

internal/duckdb/analytics_usage.go

Lines changed: 26 additions & 42 deletions
Original file line numberDiff line numberDiff line change
@@ -1662,11 +1662,8 @@ func (s *Store) GetTrendsTerms(
16621662
continue
16631663
}
16641664
messageCounts[pos]++
1665-
lower := strings.ToLower(content)
16661665
for i, term := range terms {
1667-
for _, variant := range trendVariants(term) {
1668-
counts[i][pos] += strings.Count(lower, strings.ToLower(variant))
1669-
}
1666+
counts[i][pos] += db.CountTrendOccurrences(content, term)
16701667
}
16711668
}
16721669
if err := rows.Err(); err != nil {
@@ -1677,19 +1674,6 @@ func (s *Store) GetTrendsTerms(
16771674
), nil
16781675
}
16791676

1680-
func trendVariants(term db.TrendTermInput) []string {
1681-
if len(term.Matchers) > 0 {
1682-
return term.Matchers
1683-
}
1684-
if len(term.Variants) > 0 {
1685-
return term.Variants
1686-
}
1687-
if term.Term != "" {
1688-
return []string{term.Term}
1689-
}
1690-
return nil
1691-
}
1692-
16931677
type duckRates struct {
16941678
input float64
16951679
output float64
@@ -1946,6 +1930,20 @@ func duckUsageLocalDateSQL(f db.UsageFilter) (string, any) {
19461930
func duckUsageCTE(f db.UsageFilter, sessionID string) (string, []any) {
19471931
rawSQL, args := duckUsageRawSQL(f, sessionID)
19481932
localDateSQL, localDateArg := duckUsageLocalDateSQL(f)
1933+
// Apply the local-date window BEFORE deduping so an out-of-range
1934+
// duplicate (pulled in by the padded UTC bounds) cannot win
1935+
// dedup_rank = 1 and suppress the in-range row. Mirrors the
1936+
// dedup-after-date-filter order in internal/db/usage.go.
1937+
datePred := "TRUE"
1938+
var dateArgs []any
1939+
if f.From != "" {
1940+
datePred += " AND local_date >= ?"
1941+
dateArgs = append(dateArgs, f.From)
1942+
}
1943+
if f.To != "" {
1944+
datePred += " AND local_date <= ?"
1945+
dateArgs = append(dateArgs, f.To)
1946+
}
19491947
query := fmt.Sprintf(`
19501948
WITH usage_raw AS (
19511949
%s
@@ -1976,41 +1974,33 @@ func duckUsageCTE(f db.UsageFilter, sessionID string) (string, []any) {
19761974
ELSE 'row:' || session_id || ':' || source || ':' ||
19771975
COALESCE(CAST(message_ordinal AS VARCHAR), '') || ':' ||
19781976
CAST(ts AS VARCHAR) || ':' || model
1979-
END AS dedup_group
1977+
END AS dedup_group,
1978+
%s AS local_date
19801979
FROM usage_raw
19811980
),
1981+
usage_windowed AS (
1982+
SELECT *
1983+
FROM usage_normalized
1984+
WHERE %s
1985+
),
19821986
usage_ranked AS (
19831987
SELECT *,
19841988
ROW_NUMBER() OVER (
19851989
PARTITION BY dedup_group
19861990
ORDER BY ts ASC, session_id ASC, COALESCE(message_ordinal, -1) ASC
19871991
) AS dedup_rank
1988-
FROM usage_normalized
1992+
FROM usage_windowed
19891993
),
19901994
usage_localized AS (
1991-
SELECT *,
1992-
%s AS local_date
1995+
SELECT *
19931996
FROM usage_ranked
19941997
WHERE dedup_rank = 1
1995-
)`, rawSQL, localDateSQL)
1998+
)`, rawSQL, localDateSQL, datePred)
19961999
args = append(args, localDateArg)
2000+
args = append(args, dateArgs...)
19972001
return query, args
19982002
}
19992003

2000-
func appendDuckUsageLocalDateFilter(
2001-
where string, args []any, f db.UsageFilter,
2002-
) (string, []any) {
2003-
if f.From != "" {
2004-
where += "\n\t\tAND local_date >= ?"
2005-
args = append(args, f.From)
2006-
}
2007-
if f.To != "" {
2008-
where += "\n\t\tAND local_date <= ?"
2009-
args = append(args, f.To)
2010-
}
2011-
return where, args
2012-
}
2013-
20142004
type duckUsageBucket struct {
20152005
inputTok int
20162006
outputTok int
@@ -2064,8 +2054,6 @@ func (s *Store) dailyUsageAggregateRows(
20642054
ctx context.Context, f db.UsageFilter,
20652055
) ([]duckUsageAggregateRow, error) {
20662056
cte, args := duckUsageCTE(f, "")
2067-
where := "TRUE"
2068-
where, args = appendDuckUsageLocalDateFilter(where, args, f)
20692057
query := cte + `
20702058
SELECT local_date, project, agent, model,
20712059
SUM(input_tokens_norm) AS input_tokens,
@@ -2078,7 +2066,6 @@ func (s *Store) dailyUsageAggregateRows(
20782066
SUM(CASE WHEN cost_usd IS NULL THEN cache_read_norm ELSE 0 END) AS billable_cache_read_tokens,
20792067
COALESCE(SUM(cost_usd), 0) AS explicit_cost
20802068
FROM usage_localized
2081-
WHERE ` + where + `
20822069
GROUP BY local_date, project, agent, model
20832070
ORDER BY local_date ASC, project ASC, agent ASC, model ASC`
20842071
rows, err := s.duck.QueryContext(ctx, query, args...)
@@ -2279,8 +2266,6 @@ func (s *Store) sessionUsageAggregateRows(
22792266
ctx context.Context, f db.UsageFilter, sessionID string,
22802267
) ([]duckUsageAggregateRow, error) {
22812268
cte, args := duckUsageCTE(f, sessionID)
2282-
where := "TRUE"
2283-
where, args = appendDuckUsageLocalDateFilter(where, args, f)
22842269
query := cte + `
22852270
SELECT session_id, project, agent, model,
22862271
ANY_VALUE(display_name) AS display_name,
@@ -2295,7 +2280,6 @@ func (s *Store) sessionUsageAggregateRows(
22952280
SUM(CASE WHEN cost_usd IS NULL THEN cache_read_norm ELSE 0 END) AS billable_cache_read_tokens,
22962281
COALESCE(SUM(cost_usd), 0) AS explicit_cost
22972282
FROM usage_localized
2298-
WHERE ` + where + `
22992283
GROUP BY session_id, project, agent, model
23002284
ORDER BY session_id ASC, model ASC`
23012285
rows, err := s.duck.QueryContext(ctx, query, args...)

internal/duckdb/store_contract_test.go

Lines changed: 3 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -197,10 +197,9 @@ func duckContractAnalyticsTrendsAndUsage(
197197
require.NoError(t, err)
198198
require.Equal(t, 1, tools.TotalCalls)
199199

200-
trends, err := store.GetTrendsTerms(ctx, filter, []db.TrendTermInput{{
201-
Term: "alpha",
202-
Variants: []string{"alpha"},
203-
}}, "week")
200+
trendTerms, err := db.ParseTrendTerms([]string{"alpha"})
201+
require.NoError(t, err)
202+
trends, err := store.GetTrendsTerms(ctx, filter, trendTerms, "week")
204203
require.NoError(t, err)
205204
require.Equal(t, 1, trends.Series[0].Total)
206205

internal/duckdb/store_test.go

Lines changed: 94 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -374,10 +374,9 @@ func TestStoreAnalyticsUsageAndTrends(t *testing.T) {
374374
require.NoError(t, err)
375375
assert.Equal(t, 2, signals.UnscoredSessions)
376376

377-
trends, err := store.GetTrendsTerms(ctx, filter, []db.TrendTermInput{{
378-
Term: "alpha",
379-
Variants: []string{"alpha"},
380-
}}, "week")
377+
trendTerms, err := db.ParseTrendTerms([]string{"alpha"})
378+
require.NoError(t, err)
379+
trends, err := store.GetTrendsTerms(ctx, filter, trendTerms, "week")
381380
require.NoError(t, err)
382381
assert.Equal(t, 1, trends.Series[0].Total)
383382

@@ -1391,15 +1390,14 @@ func TestTrendsTermsApplySessionFiltersAndSystemPrefixExclusion(t *testing.T) {
13911390
require.NoError(t, err)
13921391
store := NewStoreFromDB(syncer.DB())
13931392

1393+
trendTerms, err := db.ParseTrendTerms([]string{"seam"})
1394+
require.NoError(t, err)
13941395
trends, err := store.GetTrendsTerms(ctx, db.AnalyticsFilter{
13951396
From: "2026-01-22",
13961397
To: "2026-01-22",
13971398
Timezone: "UTC",
13981399
Project: "alpha",
1399-
}, []db.TrendTermInput{{
1400-
Term: "seam",
1401-
Variants: []string{"seam"},
1402-
}}, "day")
1400+
}, trendTerms, "day")
14031401
require.NoError(t, err)
14041402
require.Len(t, trends.Series, 1)
14051403
assert.Equal(t, 1, trends.Series[0].Total)
@@ -1558,6 +1556,94 @@ func TestUsageDedupesClaudeMessageIDs(t *testing.T) {
15581556
assert.Equal(t, []string{"claude-test"}, sessionUsage.Models)
15591557
}
15601558

1559+
func TestUsageDedupPrefersInRangeDuplicate(t *testing.T) {
1560+
ctx := context.Background()
1561+
local := newLocalDB(t)
1562+
require.NoError(t, local.UpsertModelPricing([]db.ModelPricing{{
1563+
ModelPattern: "claude-test",
1564+
InputPerMTok: 3,
1565+
OutputPerMTok: 15,
1566+
}}))
1567+
1568+
before := syncMessage("duck-usage-edge-a", 0, "assistant", "before midnight", "2026-01-12T23:30:00.000Z")
1569+
before.ClaudeMessageID = "edge-message"
1570+
before.ClaudeRequestID = "edge-request"
1571+
after := syncMessage("duck-usage-edge-b", 0, "assistant", "after midnight", "2026-01-13T00:30:00.000Z")
1572+
after.ClaudeMessageID = "edge-message"
1573+
after.ClaudeRequestID = "edge-request"
1574+
1575+
_, err := local.WriteSessionBatchAtomic([]db.SessionBatchWrite{
1576+
{
1577+
Session: syncSession("duck-usage-edge-a", "alpha", "edge a", "2026-01-12T23:30:00.000Z", 1),
1578+
Messages: []db.Message{before},
1579+
DataVersion: 1,
1580+
ReplaceMessages: true,
1581+
},
1582+
{
1583+
Session: syncSession("duck-usage-edge-b", "alpha", "edge b", "2026-01-13T00:30:00.000Z", 1),
1584+
Messages: []db.Message{after},
1585+
DataVersion: 1,
1586+
ReplaceMessages: true,
1587+
},
1588+
})
1589+
require.NoError(t, err)
1590+
1591+
syncer := newTestSync(t, filepath.Join(t.TempDir(), "usage-edge.duckdb"), local, SyncOptions{})
1592+
_, err = syncer.Push(ctx, true, nil)
1593+
require.NoError(t, err)
1594+
store := NewStoreFromDB(syncer.DB())
1595+
1596+
// The duplicate before midnight is outside the window but inside
1597+
// the padded UTC bounds and sorts first by timestamp. It must not
1598+
// win the dedup and suppress the in-range duplicate.
1599+
got, err := store.GetDailyUsage(ctx, db.UsageFilter{
1600+
From: "2026-01-13", To: "2026-01-13", Timezone: "UTC",
1601+
})
1602+
require.NoError(t, err)
1603+
assert.Equal(t, 1, got.Totals.InputTokens)
1604+
assert.Equal(t, 2, got.Totals.OutputTokens)
1605+
}
1606+
1607+
func TestTrendsTermsWordBoundaryAndOverlapParity(t *testing.T) {
1608+
ctx := context.Background()
1609+
local := newLocalDB(t)
1610+
start := "2026-01-22T09:00:00.000Z"
1611+
content := "seam seams seamless testing test attest"
1612+
_, err := local.WriteSessionBatchAtomic([]db.SessionBatchWrite{{
1613+
Session: syncSession("duck-trend-parity", "alpha", "trend parity", start, 1),
1614+
Messages: []db.Message{
1615+
syncMessage("duck-trend-parity", 0, "user", content, start),
1616+
},
1617+
DataVersion: 1,
1618+
ReplaceMessages: true,
1619+
}})
1620+
require.NoError(t, err)
1621+
syncer := newTestSync(t, filepath.Join(t.TempDir(), "trends-parity.duckdb"), local, SyncOptions{})
1622+
_, err = syncer.Push(ctx, true, nil)
1623+
require.NoError(t, err)
1624+
store := NewStoreFromDB(syncer.DB())
1625+
1626+
terms, err := db.ParseTrendTerms([]string{"seam", "test|testing"})
1627+
require.NoError(t, err)
1628+
filter := db.AnalyticsFilter{
1629+
From: "2026-01-22", To: "2026-01-22", Timezone: "UTC",
1630+
}
1631+
1632+
got, err := store.GetTrendsTerms(ctx, filter, terms, "day")
1633+
require.NoError(t, err)
1634+
require.Len(t, got.Series, 2)
1635+
// Word-bounded: "seamless" does not count for "seam", and
1636+
// "testing" is not double-counted via its "test" substring.
1637+
assert.Equal(t, 2, got.Series[0].Total)
1638+
assert.Equal(t, 2, got.Series[1].Total)
1639+
1640+
want, err := local.GetTrendsTerms(ctx, filter, terms, "day")
1641+
require.NoError(t, err)
1642+
require.Len(t, want.Series, 2)
1643+
assert.Equal(t, want.Series[0].Total, got.Series[0].Total)
1644+
assert.Equal(t, want.Series[1].Total, got.Series[1].Total)
1645+
}
1646+
15611647
func TestDailyUsageBreakdownsAndCacheSavings(t *testing.T) {
15621648
ctx := context.Background()
15631649
local := newLocalDB(t)

0 commit comments

Comments
 (0)