diff --git a/gateway/warmup/cache.go b/gateway/warmup/cache.go index bb4215f4..02b0fd3b 100644 --- a/gateway/warmup/cache.go +++ b/gateway/warmup/cache.go @@ -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 } @@ -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 } @@ -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, }) } } @@ -167,7 +169,7 @@ func summarize(entries []*PreCached) *Summary { continue } summary.CompletedCases++ - summary.WarmedRows += entry.Rows + summary.GroupsWritten += entry.GroupsWritten } return summary } @@ -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 { @@ -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) } @@ -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)) } } @@ -287,7 +289,7 @@ func warmupMetricOperation(aView *view.View) *gmetricx.OperationRef { warmupRunErrorKey, warmupCasesCompletedKey, warmupCasesFailedKey, - warmupRowsKey, + warmupGroupsWrittenKey, )) }) } diff --git a/gateway/warmup/cache_test.go b/gateway/warmup/cache_test.go index 94574492..208c3eb7 100644 --- a/gateway/warmup/cache_test.go +++ b/gateway/warmup/cache_test.go @@ -16,13 +16,13 @@ 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}, }, } @@ -30,12 +30,13 @@ func TestAppendPreCachedUsesEntryRows(t *testing.T) { 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) } @@ -43,7 +44,7 @@ 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"}, }, } @@ -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) { @@ -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) } @@ -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) } @@ -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) } @@ -198,7 +196,7 @@ func TestRecordWarmupViewMetrics(t *testing.T) { View: aView.Name, CompletedCases: 2, FailedCases: 1, - WarmedRows: 40, + GroupsWritten: 400, Elapsed: 1500 * time.Millisecond, }) @@ -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)) } diff --git a/go.mod b/go.mod index f8882f62..b978e3a2 100644 --- a/go.mod +++ b/go.mod @@ -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 diff --git a/go.sum b/go.sum index c456cb73..f8efa260 100644 --- a/go.sum +++ b/go.sum @@ -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= diff --git a/service/reader/service.go b/service/reader/service.go index 468e7967..68b853a7 100644 --- a/service/reader/service.go +++ b/service/reader/service.go @@ -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), @@ -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 { diff --git a/service/reader/service_metrics_test.go b/service/reader/service_metrics_test.go index 7d65633a..c8a32aa5 100644 --- a/service/reader/service_metrics_test.go +++ b/service/reader/service_metrics_test.go @@ -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 } diff --git a/warmup/cache.go b/warmup/cache.go index 46f39d8f..7f742d0e 100644 --- a/warmup/cache.go +++ b/warmup/cache.go @@ -9,7 +9,6 @@ import ( errUtils "github.com/viant/datly/shared" "github.com/viant/datly/view" "github.com/viant/sqlx/io/read/cache" - cachehash "github.com/viant/sqlx/io/read/cache/hash" "strings" "sync" "time" @@ -34,27 +33,27 @@ type ( column string label string fields string - key string } warmupEntryFn func() (*warmupEntry, error) notifierFn func() (int, *EntryResult, error) EntryResult struct { - View string - Column string - Params string - CacheKey string - FieldNames string - Elapsed string - TimeTaken time.Duration - Rows int - Error string `json:",omitempty"` + View string + Column string + Params string + WarmupKey string + MarkerKey string + FieldNames string + Elapsed string + TimeTaken time.Duration + GroupsWritten int + Error string `json:",omitempty"` } Result struct { - Rows int - Entries []*EntryResult + GroupsWritten int + Entries []*EntryResult } ) @@ -66,7 +65,7 @@ func (c *matchersCollector) populate(ctx context.Context, collector chan warmupE if err == nil { return size, nil, nil } - return size, failedEntryResult(&warmupEntry{view: c.view}, 0, 0, err), err + return size, failedEntryResult(&warmupEntry{view: c.view}, 0, err), err } }() } @@ -118,20 +117,6 @@ func (c *matchersCollector) createMetaWarmupEntry(ctx context.Context, aView *vi } return } - cacheKey, err := warmupCacheKey(cacheIndex) - if err != nil { - fmt.Printf("[INFO] cache warmup entry build error view=%s type=meta column=%s field_names=%s error=%v\n", aView.Name, input.MetaColumn, strings.Join(input.FieldNames, ","), err) - aChan <- func() (*warmupEntry, error) { - return &warmupEntry{ - view: aView, - column: input.MetaColumn, - label: input.Label, - fields: strings.Join(input.FieldNames, ","), - }, err - } - return - } - aChan <- func() (*warmupEntry, error) { return &warmupEntry{ matcher: cacheIndex, @@ -139,7 +124,6 @@ func (c *matchersCollector) createMetaWarmupEntry(ctx context.Context, aView *vi column: input.MetaColumn, label: input.Label, fields: strings.Join(input.FieldNames, ","), - key: cacheKey, }, nil } } @@ -158,19 +142,6 @@ func (c *matchersCollector) createIndexWarmupEntry(ctx context.Context, aView *v } return } - cacheKey, err := warmupCacheKey(build) - if err != nil { - fmt.Printf("[INFO] cache warmup entry build error view=%s type=index column=%s field_names=%s error=%v\n", aView.Name, cacheInput.Column, strings.Join(cacheInput.FieldNames, ","), err) - aChan <- func() (*warmupEntry, error) { - return &warmupEntry{ - view: aView, - column: cacheInput.Column, - label: cacheInput.Label, - fields: strings.Join(cacheInput.FieldNames, ","), - }, err - } - return - } build.StoredFields = view.SQLXProjectionFields(cacheInput.StoredFields) aChan <- func() (*warmupEntry, error) { @@ -180,7 +151,6 @@ func (c *matchersCollector) createIndexWarmupEntry(ctx context.Context, aView *v column: cacheInput.Column, label: cacheInput.Label, fields: strings.Join(cacheInput.FieldNames, ","), - key: cacheKey, }, nil } } @@ -239,39 +209,68 @@ func readWithChan(ctx context.Context, entry *warmupEntry, notifier chan func() func readWithErr(ctx context.Context, entry *warmupEntry) (*EntryResult, error) { started := time.Now() - fmt.Printf("[INFO] cache warmup query start start_time=%s view=%s cache=%s cache_key=%s db_connector=%s column=%s params=%s field_names=%s args=%v sql=%q\n", started.Format(time.RFC3339), entry.view.Name, cacheLabel(entry.view), entry.key, warmupConnectorLabel(entry.view), entry.column, entry.label, entry.fields, entry.matcher.Args, truncateSQL(entry.matcher.SQL)) + fmt.Printf("[INFO] cache warmup query start start_time=%s view=%s cache=%s db_connector=%s column=%s params=%s field_names=%s args=%v sql=%q\n", started.Format(time.RFC3339), entry.view.Name, cacheLabel(entry.view), warmupConnectorLabel(entry.view), entry.column, entry.label, entry.fields, entry.matcher.Args, truncateSQL(entry.matcher.SQL)) db, err := DB(entry) if err != nil { elapsed := time.Since(started) - fmt.Printf("[INFO] cache warmup query error view=%s cache_key=%s column=%s params=%s field_names=%s elapsed=%s cache_write=skipped error=%v\n", entry.view.Name, entry.key, entry.column, entry.label, entry.fields, elapsed, err) - return failedEntryResult(entry, elapsed, 0, err), err + fmt.Printf("[INFO] cache warmup query error view=%s column=%s params=%s field_names=%s elapsed=%s cache_write=skipped error=%v\n", entry.view.Name, entry.column, entry.label, entry.fields, elapsed, err) + return failedEntryResult(entry, elapsed, err), err } service, err := entry.view.Cache.Service() if err != nil { elapsed := time.Since(started) - fmt.Printf("[INFO] cache warmup query error view=%s cache_key=%s column=%s params=%s field_names=%s elapsed=%s cache_write=skipped error=%v\n", entry.view.Name, entry.key, entry.column, entry.label, entry.fields, elapsed, err) - return failedEntryResult(entry, elapsed, 0, err), err + fmt.Printf("[INFO] cache warmup query error view=%s column=%s params=%s field_names=%s elapsed=%s cache_write=skipped error=%v\n", entry.view.Name, entry.column, entry.label, entry.fields, elapsed, err) + return failedEntryResult(entry, elapsed, err), err } matcher := entry.matcher - indexed, err := service.IndexBy(indexProgressContext(ctx, entry), db, entry.column, matcher.SQL, matcher.Args, matcher) + indexResult, err := indexByWithResult(indexProgressContext(ctx, entry), service, db, entry.column, matcher.SQL, matcher.Args, matcher) elapsed := time.Since(started) + if indexResult == nil { + indexResult = &indexByResult{} + } if err != nil { - fmt.Printf("[INFO] cache warmup query error view=%s cache_key=%s column=%s params=%s field_names=%s rows=%d elapsed=%s cache_write=error error=%v\n", entry.view.Name, entry.key, entry.column, entry.label, entry.fields, indexed, elapsed, err) + fmt.Printf("[INFO] cache warmup query error view=%s warmup_key=%s marker_key=%s column=%s params=%s field_names=%s groups_written=%d elapsed=%s cache_write=error error=%v\n", entry.view.Name, indexResult.warmupKey, indexResult.markerKey, entry.column, entry.label, entry.fields, indexResult.groupsWritten, elapsed, err) indexErr := fmt.Errorf("failed to index: %w", err) - return failedEntryResult(entry, elapsed, indexed, indexErr), indexErr + result := failedEntryResult(entry, elapsed, indexErr) + result.WarmupKey = indexResult.warmupKey + result.MarkerKey = indexResult.markerKey + result.GroupsWritten = indexResult.groupsWritten + return result, indexErr + } + + fmt.Printf("[INFO] cache warmup query done view=%s cache=%s warmup_key=%s marker_key=%s db_connector=%s column=%s params=%s field_names=%s groups_written=%d elapsed=%s cache_write=success\n", entry.view.Name, cacheLabel(entry.view), indexResult.warmupKey, indexResult.markerKey, warmupConnectorLabel(entry.view), entry.column, entry.label, entry.fields, indexResult.groupsWritten, elapsed) + return &EntryResult{View: entry.view.Name, Column: entry.column, Params: entry.label, WarmupKey: indexResult.warmupKey, MarkerKey: indexResult.markerKey, FieldNames: entry.fields, Elapsed: elapsed.String(), TimeTaken: elapsed, GroupsWritten: indexResult.groupsWritten}, nil +} + +type indexByResult struct { + groupsWritten int + warmupKey string + markerKey string +} + +func indexByWithResult(ctx context.Context, service cache.Cache, db *sql.DB, column, SQL string, args []interface{}, matcher *cache.ParmetrizedQuery) (*indexByResult, error) { + if indexer, ok := service.(cache.WarmupIndexer); ok { + result, err := indexer.IndexByWithResult(ctx, db, column, SQL, args, matcher) + if result == nil { + return nil, err + } + return &indexByResult{ + groupsWritten: result.GroupsWritten, + warmupKey: result.WarmupKey, + markerKey: result.MarkerKey, + }, err } - fmt.Printf("[INFO] cache warmup query done view=%s cache=%s cache_key=%s db_connector=%s column=%s params=%s field_names=%s rows=%d elapsed=%s cache_write=success\n", entry.view.Name, cacheLabel(entry.view), entry.key, warmupConnectorLabel(entry.view), entry.column, entry.label, entry.fields, indexed, elapsed) - return &EntryResult{View: entry.view.Name, Column: entry.column, Params: entry.label, CacheKey: entry.key, FieldNames: entry.fields, Elapsed: elapsed.String(), TimeTaken: elapsed, Rows: indexed}, nil + groupsWritten, err := service.IndexBy(ctx, db, column, SQL, args, matcher) + return &indexByResult{groupsWritten: groupsWritten}, err } -func failedEntryResult(entry *warmupEntry, elapsed time.Duration, rows int, err error) *EntryResult { +func failedEntryResult(entry *warmupEntry, elapsed time.Duration, err error) *EntryResult { result := &EntryResult{ Elapsed: elapsed.String(), TimeTaken: elapsed, - Rows: rows, } if entry != nil { if entry.view != nil { @@ -279,7 +278,6 @@ func failedEntryResult(entry *warmupEntry, elapsed time.Duration, rows int, err } result.Column = entry.column result.Params = entry.label - result.CacheKey = entry.key result.FieldNames = entry.fields } if err != nil { @@ -297,24 +295,6 @@ func firstError(errors []error) error { return nil } -func warmupCacheKey(query *cache.ParmetrizedQuery) (string, error) { - if query == nil { - return "", fmt.Errorf("warmup cache key query was nil") - } - return warmupIdentityURL(query) -} - -func warmupIdentityURL(query *cache.ParmetrizedQuery) (string, error) { - if query == nil { - return "", fmt.Errorf("warmup identity query was nil") - } - SQL, _, argsMarshal, err := query.WarmupIdentity() - if err != nil { - return "", err - } - return cachehash.GenerateWithMarshal(SQL, "", "", argsMarshal) -} - func DB(entry *warmupEntry) (*sql.DB, error) { if entry.view.Cache.Warmup.Connector != nil { return entry.view.Cache.Warmup.Connector.DB() @@ -328,7 +308,7 @@ func PopulateCache(views []*view.View) (int, error) { if result == nil { return 0, err } - return result.Rows, err + return result.GroupsWritten, err } func PopulateCacheWithDetails(views []*view.View) (*Result, error) { @@ -342,7 +322,7 @@ func PopulateCacheWithDetailsContext(ctx context.Context, views []*view.View) (* result := &Result{} if len(viewsWithCache) == 0 { - fmt.Printf("[INFO] cache warmup populate done rows=0 elapsed=%s\n", time.Since(started)) + fmt.Printf("[INFO] cache warmup populate done groups_written=0 elapsed=%s\n", time.Since(started)) return result, nil } @@ -376,7 +356,7 @@ func PopulateCacheWithDetailsContext(ctx context.Context, views []*view.View) (* } if collectorSize == 0 { - fmt.Printf("[INFO] cache warmup populate done rows=0 entries=0 elapsed=%s\n", time.Since(started)) + fmt.Printf("[INFO] cache warmup populate done groups_written=0 entries=0 elapsed=%s\n", time.Since(started)) err := errUtils.CombineErrors("errors while populating cache: ", errors) if err != nil { return result, err @@ -391,7 +371,7 @@ func PopulateCacheWithDetailsContext(ctx context.Context, views []*view.View) (* entry, err := fn() if err != nil { errors = append(errors, err) - result.Entries = append(result.Entries, failedEntryResult(entry, 0, 0, err)) + result.Entries = append(result.Entries, failedEntryResult(entry, 0, err)) } else { warmupEntries = append(warmupEntries, entry) } @@ -419,7 +399,7 @@ func PopulateCacheWithDetailsContext(ctx context.Context, views []*view.View) (* entryResult, err := actual() if entryResult != nil { result.Entries = append(result.Entries, entryResult) - result.Rows += entryResult.Rows + result.GroupsWritten += entryResult.GroupsWritten } if err != nil { errors = append(errors, err) @@ -429,10 +409,10 @@ func PopulateCacheWithDetailsContext(ctx context.Context, views []*view.View) (* close(notifier) err := errUtils.CombineErrors("errors while populating cache: ", errors) if err != nil { - fmt.Printf("[INFO] cache warmup populate error rows=%d entries=%d failures=%d elapsed=%s first_error=%v\n", result.Rows, len(warmupEntries), len(errors), time.Since(started), firstError(errors)) + fmt.Printf("[INFO] cache warmup populate error groups_written=%d entries=%d failures=%d elapsed=%s first_error=%v\n", result.GroupsWritten, len(warmupEntries), len(errors), time.Since(started), firstError(errors)) return result, err } - fmt.Printf("[INFO] cache warmup populate done rows=%d entries=%d elapsed=%s\n", result.Rows, len(warmupEntries), time.Since(started)) + fmt.Printf("[INFO] cache warmup populate done groups_written=%d entries=%d elapsed=%s\n", result.GroupsWritten, len(warmupEntries), time.Since(started)) return result, nil } diff --git a/warmup/cache_test.go b/warmup/cache_test.go index b7334521..0a574ad9 100644 --- a/warmup/cache_test.go +++ b/warmup/cache_test.go @@ -17,7 +17,7 @@ import ( sqlcache "github.com/viant/sqlx/io/read/cache" ) -func TestPopulateCache(t *testing.T) { +func TestPopulateCacheWithDetails(t *testing.T) { if os.Getenv("DATLY_RUN_WARMUP_TESTS") == "" { t.Skip("set DATLY_RUN_WARMUP_TESTS=1 to run warmup integration test") } @@ -78,9 +78,10 @@ func TestPopulateCache(t *testing.T) { views = append(views, item) } - inserted, err := PopulateCache(views) + result, err := PopulateCacheWithDetails(views) assert.Nil(t, err, testCase.description) - assert.Equal(t, testCase.expectedInserted, inserted, testCase.description) + require.NotNil(t, result, testCase.description) + assert.Equal(t, testCase.expectedInserted, result.GroupsWritten, testCase.description) for _, aView := range views { cache := aView.Cache @@ -146,7 +147,7 @@ func TestWarmupWithLimitCapsConcurrency(t *testing.T) { } time.Sleep(5 * time.Millisecond) atomic.AddInt64(&active, -1) - return &EntryResult{Rows: 1}, nil + return &EntryResult{GroupsWritten: 1}, nil } notifier := make(chan func() (*EntryResult, error)) @@ -157,26 +158,13 @@ func TestWarmupWithLimitCapsConcurrency(t *testing.T) { actual := <-notifier result, err := actual() assert.Nil(t, err) - total += result.Rows + total += result.GroupsWritten } assert.Equal(t, len(entries), total) assert.LessOrEqual(t, atomic.LoadInt64(&maxActive), int64(maxWarmupConcurrency)) } -func TestWarmupCacheKeyNormalizesNilArgs(t *testing.T) { - nilArgsKey, err := warmupCacheKey(&sqlcache.ParmetrizedQuery{SQL: "SELECT * FROM events", Args: nil}) - assert.Nil(t, err) - - emptyArgsKey, err := warmupCacheKey(&sqlcache.ParmetrizedQuery{SQL: "SELECT * FROM events", Args: []interface{}{}}) - assert.Nil(t, err) - - assert.Equal(t, emptyArgsKey, nilArgsKey) - - _, err = warmupCacheKey(nil) - assert.ErrorContains(t, err, "query was nil") -} - func TestIndexProgressContext(t *testing.T) { aView := &view.View{ Name: "performanceTimeline", @@ -228,7 +216,7 @@ func TestIndexProgressContext(t *testing.T) { require.True(t, actual.Done) } -func TestWarmupFieldNamesAffectGeneratedCacheKey(t *testing.T) { +func TestWarmupFieldNamesAffectGeneratedProjection(t *testing.T) { resourcePath := path.Join(t.TempDir(), "resource.yaml") require.NoError(t, os.WriteFile(resourcePath, []byte(` CacheProviders: @@ -276,8 +264,6 @@ Views: builder := reader.NewBuilder() fullQuery, err := builder.CacheSQL(context.Background(), aView, input[0].Selector) require.NoError(t, err) - fullKey, err := warmupCacheKey(fullQuery) - require.NoError(t, err) aView.Cache.Warmup.FieldNames = []string{"Quantity"} fieldInput, err := aView.Cache.GenerateCacheInput(context.Background()) @@ -291,11 +277,8 @@ Views: fieldQuery, err := builder.CacheSQL(context.Background(), aView, fieldInput[0].Selector) require.NoError(t, err) - fieldKey, err := warmupCacheKey(fieldQuery) - require.NoError(t, err) assert.NotEqual(t, fullQuery.SQL, fieldQuery.SQL) - assert.NotEqual(t, fullKey, fieldKey) assert.Contains(t, fieldQuery.SQL, "quantity") }