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
68 changes: 35 additions & 33 deletions gateway/warmup/cache.go
Original file line number Diff line number Diff line change
Expand Up @@ -23,36 +23,37 @@ const (
warmupRunErrorKey = "run.error"
warmupCasesCompletedKey = "cases.completed"
warmupCasesFailedKey = "cases.failed"
warmupRowsKey = "rows"
warmupGroupsWrittenKey = "groupsWritten"
warmupMetricFallbackPkg = "datly"
warmupMetricRecentBuckets = 2
)

type PreCachables func(ctx context.Context, method, matchingURI string) ([]*view.View, error)
type PreCached struct {
URI string
View string
Column string
Params string
CacheKey string
FieldNames string `json:",omitempty"`
Elapsed string
TimeTaken time.Duration
Rows int
Error string `json:"error,omitempty"`
URI string
View string
Column string
Params string
WarmupKey string
MarkerKey string `json:",omitempty"`
FieldNames string `json:",omitempty"`
Elapsed string
TimeTaken time.Duration
GroupsWritten int `json:"groupsWritten,omitempty"`
Error string `json:"error,omitempty"`
}

type Summary struct {
CompletedCases int `json:"completedCases"`
FailedCases int `json:"failedCases"`
WarmedRows int `json:"warmedRows"`
GroupsWritten int `json:"groupsWritten,omitempty"`
}

type viewSummary struct {
View string
CompletedCases int
FailedCases int
WarmedRows int
GroupsWritten int
Elapsed time.Duration
}

Expand Down Expand Up @@ -106,11 +107,11 @@ func PreCache(ctx context.Context, lookup PreCachables, warmupURIs ...string) *R
setErr(e)
}
elapsed := time.Now().Sub(startTime)
rows := 0
groupsWritten := 0
if result != nil {
rows = result.Rows
groupsWritten = result.GroupsWritten
}
fmt.Printf("[INFO] cache warmup uri done uri=%s rows=%d elapsed=%s\n", URI, rows, elapsed)
fmt.Printf("[INFO] cache warmup uri done uri=%s groups_written=%d elapsed=%s\n", URI, groupsWritten, elapsed)
if result == nil {
return
}
Expand Down Expand Up @@ -139,16 +140,17 @@ func appendPreCached(response *Response, URI string, result *warmup.Result) {
continue
}
response.PreCached = append(response.PreCached, &PreCached{
URI: URI,
View: entry.View,
Column: entry.Column,
Params: entry.Params,
CacheKey: entry.CacheKey,
FieldNames: entry.FieldNames,
Elapsed: entry.Elapsed,
TimeTaken: entry.TimeTaken,
Rows: entry.Rows,
Error: entry.Error,
URI: URI,
View: entry.View,
Column: entry.Column,
Params: entry.Params,
WarmupKey: entry.WarmupKey,
MarkerKey: entry.MarkerKey,
FieldNames: entry.FieldNames,
Elapsed: entry.Elapsed,
TimeTaken: entry.TimeTaken,
GroupsWritten: entry.GroupsWritten,
Error: entry.Error,
})
}
}
Expand All @@ -167,7 +169,7 @@ func summarize(entries []*PreCached) *Summary {
continue
}
summary.CompletedCases++
summary.WarmedRows += entry.Rows
summary.GroupsWritten += entry.GroupsWritten
}
return summary
}
Expand All @@ -192,7 +194,7 @@ func summarizeByView(entries []*warmup.EntryResult) []*viewSummary {
continue
}
current.CompletedCases++
current.WarmedRows += entry.Rows
current.GroupsWritten += entry.GroupsWritten
}
result := make([]*viewSummary, 0, len(index))
for _, item := range index {
Expand All @@ -210,12 +212,12 @@ func logViewSummaries(uri string, views []*view.View, result *warmup.Result) {
}
viewsIndex := indexViewsByName(views)
for _, summary := range summarizeByView(result.Entries) {
fmt.Printf("[INFO] cache warmup view summary uri=%s view=%s completed_cases=%d failed_cases=%d warmed_rows=%d elapsed=%s\n",
fmt.Printf("[INFO] cache warmup view summary uri=%s view=%s completed_cases=%d failed_cases=%d groups_written=%d elapsed=%s\n",
uri,
summary.View,
summary.CompletedCases,
summary.FailedCases,
summary.WarmedRows,
summary.GroupsWritten,
summary.Elapsed)
recordWarmupViewMetrics(viewsIndex[summary.View], summary)
}
Expand Down Expand Up @@ -265,8 +267,8 @@ func recordWarmupViewMetrics(aView *view.View, summary *viewSummary) {
if summary.FailedCases > 0 {
operation.IncrementValueBy(warmupCasesFailedKey, int64(summary.FailedCases))
}
if summary.WarmedRows > 0 {
operation.IncrementValueBy(warmupRowsKey, int64(summary.WarmedRows))
if summary.GroupsWritten > 0 {
operation.IncrementValueBy(warmupGroupsWrittenKey, int64(summary.GroupsWritten))
}
}

Expand All @@ -287,7 +289,7 @@ func warmupMetricOperation(aView *view.View) *gmetricx.OperationRef {
warmupRunErrorKey,
warmupCasesCompletedKey,
warmupCasesFailedKey,
warmupRowsKey,
warmupGroupsWrittenKey,
))
})
}
Expand Down
48 changes: 23 additions & 25 deletions gateway/warmup/cache_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,34 +16,35 @@ import (
"github.com/viant/gmetric/stat"
)

func TestAppendPreCachedUsesEntryRows(t *testing.T) {
func TestAppendPreCachedUsesGroupsWritten(t *testing.T) {
response := &Response{}
result := &datlywarmup.Result{
Rows: 30,
GroupsWritten: 300,
Entries: []*datlywarmup.EntryResult{
{View: "periodSummary#", Column: "order_id", Params: "Period=today", CacheKey: "cache://today", Elapsed: "1s", TimeTaken: time.Second, Rows: 10},
{View: "periodSummary#", Column: "order_id", Params: "Period=month", CacheKey: "cache://month", FieldNames: "OrderId,Spend", Elapsed: "2s", TimeTaken: 2 * time.Second, Rows: 20},
{View: "periodSummary#", Column: "order_id", Params: "Period=today", WarmupKey: "cache://today", Elapsed: "1s", TimeTaken: time.Second, GroupsWritten: 100},
{View: "periodSummary#", Column: "order_id", Params: "Period=month", WarmupKey: "cache://month", MarkerKey: "order_id#cache://month", FieldNames: "OrderId,Spend", Elapsed: "2s", TimeTaken: 2 * time.Second, GroupsWritten: 200},
},
}

appendPreCached(response, "/v1/api/cache/warmup/order", result)

require.Len(t, response.PreCached, 2)
require.Equal(t, "Period=today", response.PreCached[0].Params)
require.Equal(t, "cache://today", response.PreCached[0].CacheKey)
require.Equal(t, 10, response.PreCached[0].Rows)
require.Equal(t, "cache://today", response.PreCached[0].WarmupKey)
require.Equal(t, 100, response.PreCached[0].GroupsWritten)
require.Equal(t, "Period=month", response.PreCached[1].Params)
require.Equal(t, "cache://month", response.PreCached[1].CacheKey)
require.Equal(t, "cache://month", response.PreCached[1].WarmupKey)
require.Equal(t, "order_id#cache://month", response.PreCached[1].MarkerKey)
require.Equal(t, "OrderId,Spend", response.PreCached[1].FieldNames)
require.Equal(t, 20, response.PreCached[1].Rows)
require.Equal(t, 200, response.PreCached[1].GroupsWritten)
require.Equal(t, "/v1/api/cache/warmup/order", response.PreCached[1].URI)
}

func TestAppendPreCachedPreservesEntryErrors(t *testing.T) {
response := &Response{}
result := &datlywarmup.Result{
Entries: []*datlywarmup.EntryResult{
{View: "diagnostics", Column: "ad_order_id", Params: "From=2026-07-02", CacheKey: "cache://today", Elapsed: "250ms", TimeTaken: 250 * time.Millisecond, Rows: 7, Error: "failed to index"},
{View: "diagnostics", Column: "ad_order_id", Params: "From=2026-07-02", WarmupKey: "cache://today", Elapsed: "250ms", TimeTaken: 250 * time.Millisecond, GroupsWritten: 7, Error: "failed to index"},
},
}

Expand All @@ -55,16 +56,16 @@ func TestAppendPreCachedPreservesEntryErrors(t *testing.T) {

func TestSummarize(t *testing.T) {
summary := summarize([]*PreCached{
{Rows: 10, TimeTaken: 100 * time.Millisecond},
{Rows: 20, TimeTaken: 200 * time.Millisecond},
{Rows: 99, TimeTaken: 300 * time.Millisecond, Error: "failed to index"},
{Rows: 30, TimeTaken: 400 * time.Millisecond},
{GroupsWritten: 100, TimeTaken: 100 * time.Millisecond},
{GroupsWritten: 200, TimeTaken: 200 * time.Millisecond},
{GroupsWritten: 999, TimeTaken: 300 * time.Millisecond, Error: "failed to index"},
{GroupsWritten: 300, TimeTaken: 400 * time.Millisecond},
})

require.NotNil(t, summary)
require.Equal(t, 3, summary.CompletedCases)
require.Equal(t, 1, summary.FailedCases)
require.Equal(t, 60, summary.WarmedRows)
require.Equal(t, 600, summary.GroupsWritten)
}

func TestSummarizeEmpty(t *testing.T) {
Expand All @@ -73,27 +74,26 @@ func TestSummarizeEmpty(t *testing.T) {
require.NotNil(t, summary)
require.Zero(t, summary.CompletedCases)
require.Zero(t, summary.FailedCases)
require.Zero(t, summary.WarmedRows)
}

func TestSummarizeByView(t *testing.T) {
summaries := summarizeByView([]*datlywarmup.EntryResult{
{View: "periodSummary#", Rows: 10, TimeTaken: time.Second},
{View: "periodSummary#", Rows: 99, TimeTaken: 2 * time.Second, Error: "failed"},
{View: "timeline#", Rows: 20, TimeTaken: 3 * time.Second},
{View: "periodSummary#", Rows: 30, TimeTaken: 4 * time.Second},
{View: "periodSummary#", GroupsWritten: 100, TimeTaken: time.Second},
{View: "periodSummary#", GroupsWritten: 999, TimeTaken: 2 * time.Second, Error: "failed"},
{View: "timeline#", GroupsWritten: 200, TimeTaken: 3 * time.Second},
{View: "periodSummary#", GroupsWritten: 300, TimeTaken: 4 * time.Second},
})

require.Len(t, summaries, 2)
require.Equal(t, "periodSummary#", summaries[0].View)
require.Equal(t, 2, summaries[0].CompletedCases)
require.Equal(t, 1, summaries[0].FailedCases)
require.Equal(t, 40, summaries[0].WarmedRows)
require.Equal(t, 400, summaries[0].GroupsWritten)
require.Equal(t, 7*time.Second, summaries[0].Elapsed)
require.Equal(t, "timeline#", summaries[1].View)
require.Equal(t, 1, summaries[1].CompletedCases)
require.Equal(t, 0, summaries[1].FailedCases)
require.Equal(t, 20, summaries[1].WarmedRows)
require.Equal(t, 200, summaries[1].GroupsWritten)
require.Equal(t, 3*time.Second, summaries[1].Elapsed)
}

Expand Down Expand Up @@ -149,7 +149,6 @@ Connectors:
require.NotNil(t, response.Summary)
require.Equal(t, 0, response.Summary.CompletedCases)
require.Equal(t, 1, response.Summary.FailedCases)
require.Zero(t, response.Summary.WarmedRows)
require.Len(t, response.PreCached, 1)
}

Expand All @@ -162,7 +161,6 @@ func TestPreCacheLookupFailureAccounting(t *testing.T) {
require.NotNil(t, response.Summary)
require.Equal(t, 0, response.Summary.CompletedCases)
require.Equal(t, 1, response.Summary.FailedCases)
require.Zero(t, response.Summary.WarmedRows)
require.Len(t, response.PreCached, 1)
require.Equal(t, "lookup failed", response.PreCached[0].Error)
}
Expand Down Expand Up @@ -198,7 +196,7 @@ func TestRecordWarmupViewMetrics(t *testing.T) {
View: aView.Name,
CompletedCases: 2,
FailedCases: 1,
WarmedRows: 40,
GroupsWritten: 400,
Elapsed: 1500 * time.Millisecond,
})

Expand All @@ -207,6 +205,6 @@ func TestRecordWarmupViewMetrics(t *testing.T) {
require.Equal(t, int64(1), metrics.LookupOperationCumulativeMetric(metricName, warmupRunErrorKey))
require.Equal(t, int64(2), metrics.LookupOperationCumulativeMetric(metricName, warmupCasesCompletedKey))
require.Equal(t, int64(1), metrics.LookupOperationCumulativeMetric(metricName, warmupCasesFailedKey))
require.Equal(t, int64(40), metrics.LookupOperationCumulativeMetric(metricName, warmupRowsKey))
require.Equal(t, int64(400), metrics.LookupOperationCumulativeMetric(metricName, warmupGroupsWrittenKey))
require.GreaterOrEqual(t, metrics.LookupOperationCumulativeMetric(metricName, stat.CounterTimeTakenKey), int64(1500))
}
2 changes: 1 addition & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@ require (
github.com/viant/parsly v0.3.3
github.com/viant/pgo v0.11.0
github.com/viant/scy v0.33.1
github.com/viant/sqlx v0.23.1-0.20260803165008-da07533d2e8f
github.com/viant/sqlx v0.23.1-0.20260807211629-027861517984
github.com/viant/structql v0.5.4
github.com/viant/toolbox v0.37.0
github.com/viant/velty v0.4.1-0.20260408224432-5a1c31e1bd87
Expand Down
2 changes: 2 additions & 0 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -1196,6 +1196,8 @@ github.com/viant/sqlparser v0.12.1-0.20260409013525-147f8fc299b7 h1:2FdVturjHBSQ
github.com/viant/sqlparser v0.12.1-0.20260409013525-147f8fc299b7/go.mod h1:2QRGiGZYk2/pjhORGG1zLVQ9JO+bXFhqIVi31mkCRPg=
github.com/viant/sqlx v0.23.1-0.20260803165008-da07533d2e8f h1:yte+MMDo1mWS6+YpM1OYmfZe3+ieYZXsCWfjAf4c+PY=
github.com/viant/sqlx v0.23.1-0.20260803165008-da07533d2e8f/go.mod h1:yZOQRVCMZAkexsTaoqCPGJvsNO2qajQRU1VuYu23fX8=
github.com/viant/sqlx v0.23.1-0.20260807211629-027861517984 h1:/ayViIofvv1pA8wrN3ZhPmJZb0kFpfycqC6dMWXPOjg=
github.com/viant/sqlx v0.23.1-0.20260807211629-027861517984/go.mod h1:yZOQRVCMZAkexsTaoqCPGJvsNO2qajQRU1VuYu23fX8=
github.com/viant/structology v0.9.0 h1:ibR/XmdQ3+/4XW3JK+pXRqugSnxJOm2bmIvbQ0hqztY=
github.com/viant/structology v0.9.0/go.mod h1:AAFeViwniqua61sTKdOz/zlbLpN5vE4OVhDoiZJaMgA=
github.com/viant/structql v0.5.4 h1:bMdcOpzU8UMoe5OBcyJVRxLAndvU1oj3ysvPUgBckCI=
Expand Down
19 changes: 17 additions & 2 deletions service/reader/service.go
Original file line number Diff line number Diff line change
Expand Up @@ -1054,7 +1054,7 @@ func logCacheRead(ctx context.Context, aView *view.View, stats *cache.Stats, ela
return
}
recordCacheReadMetrics(aView, stats)
fmt.Printf("[INFO] datly cache read reqTraceId=%s view=%s source=%s type=%s found_warmup=%t found_lazy=%t records=%d rows=%d namespace=%s set=%s elapsed=%s args=%v\n",
fmt.Printf("[INFO] datly cache read reqTraceId=%s view=%s source=%s type=%s found_warmup=%t found_lazy=%t records=%d rows=%d namespace=%s set=%s elapsed=%s args=%v%s\n",
reqTraceID(ctx),
aView.Name,
cacheReadSource(stats),
Expand All @@ -1066,7 +1066,22 @@ func logCacheRead(ctx context.Context, aView *view.View, stats *cache.Stats, ela
stats.Namespace,
stats.Dataset,
elapsed,
args)
args,
warmupReadKeysSuffix(stats))
}

func warmupReadKeysSuffix(stats *cache.Stats) string {
if stats == nil {
return ""
}
var result string
if stats.WarmupKey != "" {
result += " warmup_key=" + stats.WarmupKey
}
if stats.MarkerKey != "" {
result += " marker_key=" + stats.MarkerKey
}
return result
}

func reqTraceID(ctx context.Context) string {
Expand Down
39 changes: 39 additions & 0 deletions service/reader/service_metrics_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -99,6 +99,45 @@ func TestRecordCacheReadMetrics(t *testing.T) {
}
}

func TestWarmupReadKeysSuffix(t *testing.T) {
testCases := []struct {
description string
stats *cache.Stats
expected string
}{
{
description: "nil stats",
expected: "",
},
{
description: "empty stats",
stats: &cache.Stats{},
expected: "",
},
{
description: "warmup key only",
stats: &cache.Stats{WarmupKey: "warmup-123"},
expected: " warmup_key=warmup-123",
},
{
description: "marker key only",
stats: &cache.Stats{MarkerKey: "order_id#warmup-123"},
expected: " marker_key=order_id#warmup-123",
},
{
description: "warmup and marker keys",
stats: &cache.Stats{WarmupKey: "warmup-123", MarkerKey: "order_id#warmup-123"},
expected: " warmup_key=warmup-123 marker_key=order_id#warmup-123",
},
}

for _, testCase := range testCases {
t.Run(testCase.description, func(t *testing.T) {
require.Equal(t, testCase.expected, warmupReadKeysSuffix(testCase.stats))
})
}
}

type metricsTestRow struct {
ID int
}
Expand Down
Loading
Loading