From 47764345f66ead5b804bcaddcd3abc4e47ecc036 Mon Sep 17 00:00:00 2001 From: Denozordec Date: Thu, 9 Apr 2026 16:04:53 +0700 Subject: [PATCH] feat: enhance CDN source management with last refreshed timestamp Added functionality to track the last refreshed timestamp for CDN sources. Updated the database schema and relevant methods to include the last refreshed timestamp during creation and updates. Implemented logic to skip fetching CDN sources based on their refresh interval, improving efficiency in the module prefix collection process. Enhanced the data retrieval methods to support the new timestamp field, ensuring accurate state management for CDN sources. --- internal/pipeline/refresh.go | 52 ++++++++++++++++++- internal/repository/postgres.go | 12 ++++- internal/repository/postgres_entities.go | 32 ++++++++---- internal/store/backend.go | 30 ++++++----- internal/store/memory_crud.go | 5 ++ ...0007_cdn_source_last_refreshed_at.down.sql | 2 + ...000007_cdn_source_last_refreshed_at.up.sql | 2 + ...0007_cdn_source_last_refreshed_at.down.sql | 2 + ...000007_cdn_source_last_refreshed_at.up.sql | 2 + 9 files changed, 112 insertions(+), 27 deletions(-) create mode 100644 migrations/postgres/000007_cdn_source_last_refreshed_at.down.sql create mode 100644 migrations/postgres/000007_cdn_source_last_refreshed_at.up.sql create mode 100644 migrations/sqlite/000007_cdn_source_last_refreshed_at.down.sql create mode 100644 migrations/sqlite/000007_cdn_source_last_refreshed_at.up.sql diff --git a/internal/pipeline/refresh.go b/internal/pipeline/refresh.go index e8226e2..d8eda6c 100644 --- a/internal/pipeline/refresh.go +++ b/internal/pipeline/refresh.go @@ -200,7 +200,16 @@ func collectModulePrefixRows(ctx context.Context, st store.Backend, hc *http.Cli return nil, err } var rows []store.PrefixRow + latestCDNRows := latestCDNRowsBySource(st, tenantID) for _, src := range sources { + sourceKey := "cdn:" + src.ID + now := time.Now().UTC() + if shouldSkipCDNSourceFetch(src, now) { + if cached := latestCDNRows[sourceKey]; len(cached) > 0 { + rows = append(rows, cached...) + continue + } + } u := strings.TrimSpace(src.URL) if u == "" { continue @@ -224,10 +233,14 @@ func collectModulePrefixRows(ctx context.Context, st store.Backend, hc *http.Cli return nil, err } etag := strings.TrimSpace(resp.Header.Get("ETag")) + patch := &store.CDNSourcePatch{} if etag != "" && etag != strings.TrimSpace(src.Etag) { e := etag - _, _ = st.UpdateCDNSource(tenantID, moduleID, src.ID, &store.CDNSourcePatch{Etag: &e}) + patch.Etag = &e } + refreshedAt := now + patch.LastRefreshedAt = &refreshedAt + _, _ = st.UpdateCDNSource(tenantID, moduleID, src.ID, patch) pfxs, err := ExtractCIDRs(string(body), src.SourceKind, src.PrefixPath) if err != nil { return nil, fmt.Errorf("cdn parse %s: %w", u, err) @@ -238,7 +251,7 @@ func collectModulePrefixRows(ctx context.Context, st store.Backend, hc *http.Cli c := *mod.DefaultCommunityID comm = &c } - rows = append(rows, store.PrefixRow{Prefix: pfx.String(), CommunityID: comm, Source: "cdn:" + src.ID}) + rows = append(rows, store.PrefixRow{Prefix: pfx.String(), CommunityID: comm, Source: sourceKey}) } } return rows, nil @@ -293,6 +306,41 @@ func collectModulePrefixRows(ctx context.Context, st store.Backend, hc *http.Cli } } +func shouldSkipCDNSourceFetch(src *store.CDNSource, now time.Time) bool { + if src == nil || src.RefreshIntervalSec == nil || *src.RefreshIntervalSec <= 0 || src.LastRefreshedAt == nil { + return false + } + nextRefreshAt := src.LastRefreshedAt.UTC().Add(time.Duration(*src.RefreshIntervalSec) * time.Second) + return now.UTC().Before(nextRefreshAt) +} + +func latestCDNRowsBySource(st store.Backend, tenantID string) map[string][]store.PrefixRow { + out := make(map[string][]store.PrefixRow) + if st == nil { + return out + } + revs, _, _ := st.ListRevisions(tenantID, "", "", 1) + if len(revs) == 0 || strings.TrimSpace(revs[0].ID) == "" { + return out + } + revID := strings.TrimSpace(revs[0].ID) + cursor := "" + for { + page, next, more := st.ListRevisionPrefixes(tenantID, revID, cursor, 2000) + for _, row := range page { + if !strings.HasPrefix(strings.TrimSpace(row.Source), "cdn:") { + continue + } + out[row.Source] = append(out[row.Source], row) + } + if !more || strings.TrimSpace(next) == "" { + break + } + cursor = next + } + return out +} + type dohJSONAnswer struct { Type int `json:"type"` Data string `json:"data"` diff --git a/internal/repository/postgres.go b/internal/repository/postgres.go index 65b1e3b..c7c5459 100644 --- a/internal/repository/postgres.go +++ b/internal/repository/postgres.go @@ -698,7 +698,7 @@ func (p *Postgres) CreateRollbackRevision(tenantID, sourceRevisionID string) (st newID := uuid.NewString() parent := sourceRevisionID meta, _ := json.Marshal(map[string]any{ - "preview_fragments": src.PreviewFragments, + "preview_fragments": src.PreviewFragments, "materialized_prefix_count": src.MaterializedPrefixCount, }) var modArg any @@ -883,7 +883,7 @@ func (p *Postgres) CreateRenderRevision(revisionID, tenantID, moduleID string, p previewFragments = map[string]string{} } meta, err := json.Marshal(map[string]any{ - "preview_fragments": previewFragments, + "preview_fragments": previewFragments, "materialized_prefix_count": len(prefixes), }) if err != nil { @@ -1178,6 +1178,14 @@ func nullInt32Ptr(i *int) *int32 { return &v } +func nullTimePtr(t *time.Time) *time.Time { + if t == nil { + return nil + } + v := t.UTC() + return &v +} + func nullJSON(s string) *string { if strings.TrimSpace(s) == "" { v := "{}" diff --git a/internal/repository/postgres_entities.go b/internal/repository/postgres_entities.go index 56b892e..bc014f4 100644 --- a/internal/repository/postgres_entities.go +++ b/internal/repository/postgres_entities.go @@ -23,7 +23,7 @@ func (p *Postgres) ListCDNSources(tenantID, moduleID string) ([]*store.CDNSource } ctx := context.Background() rows, err := p.pool.Query(ctx, ` - SELECT id::text, source_kind, url, COALESCE(prefix_path,''), COALESCE(etag,''), refresh_interval_sec, community_id::text + SELECT id::text, source_kind, url, COALESCE(prefix_path,''), COALESCE(etag,''), refresh_interval_sec, community_id::text, last_refreshed_at FROM module_cdn_source WHERE module_id=$1 ORDER BY url`, moduleID) if err != nil { return nil, err @@ -35,7 +35,8 @@ func (p *Postgres) ListCDNSources(tenantID, moduleID string) ([]*store.CDNSource s.ModuleID = moduleID var ri *int32 var comm *string - if err := rows.Scan(&s.ID, &s.SourceKind, &s.URL, &s.PrefixPath, &s.Etag, &ri, &comm); err != nil { + var last *time.Time + if err := rows.Scan(&s.ID, &s.SourceKind, &s.URL, &s.PrefixPath, &s.Etag, &ri, &comm, &last); err != nil { continue } if ri != nil { @@ -43,6 +44,10 @@ func (p *Postgres) ListCDNSources(tenantID, moduleID string) ([]*store.CDNSource s.RefreshIntervalSec = &v } s.CommunityID = strOrNil(comm) + if last != nil { + t := last.UTC() + s.LastRefreshedAt = &t + } out = append(out, &s) } return out, nil @@ -62,9 +67,9 @@ func (p *Postgres) CreateCDNSource(tenantID, moduleID string, in *store.CDNSourc ctx := context.Background() id := uuid.NewString() _, err = p.pool.Exec(ctx, ` - INSERT INTO module_cdn_source (id, module_id, source_kind, url, prefix_path, etag, refresh_interval_sec, community_id) - VALUES ($1,$2,$3,$4,$5,$6,$7, NULLIF($8::uuid, '00000000-0000-0000-0000-000000000000'::uuid))`, - id, moduleID, in.SourceKind, strings.TrimSpace(in.URL), strings.TrimSpace(in.PrefixPath), in.Etag, nullInt32Ptr(in.RefreshIntervalSec), uuidOrNilPtr(in.CommunityID)) + INSERT INTO module_cdn_source (id, module_id, source_kind, url, prefix_path, etag, refresh_interval_sec, community_id, last_refreshed_at) + VALUES ($1,$2,$3,$4,$5,$6,$7, NULLIF($8::uuid, '00000000-0000-0000-0000-000000000000'::uuid), $9)`, + id, moduleID, in.SourceKind, strings.TrimSpace(in.URL), strings.TrimSpace(in.PrefixPath), in.Etag, nullInt32Ptr(in.RefreshIntervalSec), uuidOrNilPtr(in.CommunityID), nullTimePtr(in.LastRefreshedAt)) if err != nil { return nil, err } @@ -76,9 +81,10 @@ func (p *Postgres) getCDNSource(ctx context.Context, moduleID, id string) (*stor s.ModuleID = moduleID var ri *int32 var comm *string + var last *time.Time err := p.pool.QueryRow(ctx, ` - SELECT id::text, source_kind, url, COALESCE(prefix_path,''), COALESCE(etag,''), refresh_interval_sec, community_id::text - FROM module_cdn_source WHERE id=$1 AND module_id=$2`, id, moduleID).Scan(&s.ID, &s.SourceKind, &s.URL, &s.PrefixPath, &s.Etag, &ri, &comm) + SELECT id::text, source_kind, url, COALESCE(prefix_path,''), COALESCE(etag,''), refresh_interval_sec, community_id::text, last_refreshed_at + FROM module_cdn_source WHERE id=$1 AND module_id=$2`, id, moduleID).Scan(&s.ID, &s.SourceKind, &s.URL, &s.PrefixPath, &s.Etag, &ri, &comm, &last) if err != nil { return nil, err } @@ -87,6 +93,10 @@ func (p *Postgres) getCDNSource(ctx context.Context, moduleID, id string) (*stor s.RefreshIntervalSec = &v } s.CommunityID = strOrNil(comm) + if last != nil { + t := last.UTC() + s.LastRefreshedAt = &t + } return &s, nil } @@ -131,12 +141,16 @@ func (p *Postgres) UpdateCDNSource(tenantID, moduleID, sourceID string, patch *s cur.CommunityID = &v } } + if patch.LastRefreshedAt != nil { + t := patch.LastRefreshedAt.UTC() + cur.LastRefreshedAt = &t + } ctx := context.Background() _, err = p.pool.Exec(ctx, ` UPDATE module_cdn_source SET source_kind=$3, url=$4, prefix_path=$5, etag=$6, refresh_interval_sec=$7, - community_id=NULLIF($8::uuid, '00000000-0000-0000-0000-000000000000'::uuid), updated_at=now() + community_id=NULLIF($8::uuid, '00000000-0000-0000-0000-000000000000'::uuid), last_refreshed_at=$9, updated_at=now() WHERE id=$1 AND module_id=$2`, - sourceID, moduleID, cur.SourceKind, cur.URL, cur.PrefixPath, cur.Etag, nullInt32Ptr(cur.RefreshIntervalSec), uuidOrNilPtr(cur.CommunityID)) + sourceID, moduleID, cur.SourceKind, cur.URL, cur.PrefixPath, cur.Etag, nullInt32Ptr(cur.RefreshIntervalSec), uuidOrNilPtr(cur.CommunityID), nullTimePtr(cur.LastRefreshedAt)) if err != nil { return nil, err } diff --git a/internal/store/backend.go b/internal/store/backend.go index 1d8a95a..fcd02fe 100644 --- a/internal/store/backend.go +++ b/internal/store/backend.go @@ -94,23 +94,25 @@ type ModulePatch struct { // CDNSource is a row under a CDN module. type CDNSource struct { - ID string `json:"id,omitempty"` - ModuleID string `json:"module_id,omitempty"` - SourceKind string `json:"source_kind"` - URL string `json:"url"` - PrefixPath string `json:"prefix_path,omitempty"` - Etag string `json:"etag"` - RefreshIntervalSec *int `json:"refresh_interval_sec"` - CommunityID *string `json:"community_id"` + ID string `json:"id,omitempty"` + ModuleID string `json:"module_id,omitempty"` + SourceKind string `json:"source_kind"` + URL string `json:"url"` + PrefixPath string `json:"prefix_path,omitempty"` + Etag string `json:"etag"` + RefreshIntervalSec *int `json:"refresh_interval_sec"` + CommunityID *string `json:"community_id"` + LastRefreshedAt *time.Time `json:"-"` } type CDNSourcePatch struct { - SourceKind *string `json:"source_kind,omitempty"` - URL *string `json:"url,omitempty"` - PrefixPath *string `json:"prefix_path,omitempty"` - Etag *string `json:"etag,omitempty"` - RefreshIntervalSec *int `json:"refresh_interval_sec,omitempty"` - CommunityID *string `json:"community_id,omitempty"` + SourceKind *string `json:"source_kind,omitempty"` + URL *string `json:"url,omitempty"` + PrefixPath *string `json:"prefix_path,omitempty"` + Etag *string `json:"etag,omitempty"` + RefreshIntervalSec *int `json:"refresh_interval_sec,omitempty"` + CommunityID *string `json:"community_id,omitempty"` + LastRefreshedAt *time.Time `json:"-"` } type ASEntry struct { diff --git a/internal/store/memory_crud.go b/internal/store/memory_crud.go index 93168d9..56d6ca8 100644 --- a/internal/store/memory_crud.go +++ b/internal/store/memory_crud.go @@ -143,6 +143,7 @@ func (m *Memory) CreateCDNSource(tenantID, moduleID string, in *CDNSource) (*CDN Etag: in.Etag, RefreshIntervalSec: in.RefreshIntervalSec, CommunityID: in.CommunityID, + LastRefreshedAt: in.LastRefreshedAt, } m.cdnSources[id] = s return s, nil @@ -184,6 +185,10 @@ func (m *Memory) UpdateCDNSource(tenantID, moduleID, sourceID string, patch *CDN s.CommunityID = &v } } + if patch.LastRefreshedAt != nil { + t := patch.LastRefreshedAt.UTC() + s.LastRefreshedAt = &t + } return s, nil } diff --git a/migrations/postgres/000007_cdn_source_last_refreshed_at.down.sql b/migrations/postgres/000007_cdn_source_last_refreshed_at.down.sql new file mode 100644 index 0000000..33d0552 --- /dev/null +++ b/migrations/postgres/000007_cdn_source_last_refreshed_at.down.sql @@ -0,0 +1,2 @@ +ALTER TABLE module_cdn_source +DROP COLUMN last_refreshed_at; diff --git a/migrations/postgres/000007_cdn_source_last_refreshed_at.up.sql b/migrations/postgres/000007_cdn_source_last_refreshed_at.up.sql new file mode 100644 index 0000000..d54d249 --- /dev/null +++ b/migrations/postgres/000007_cdn_source_last_refreshed_at.up.sql @@ -0,0 +1,2 @@ +ALTER TABLE module_cdn_source +ADD COLUMN last_refreshed_at TIMESTAMPTZ NULL; diff --git a/migrations/sqlite/000007_cdn_source_last_refreshed_at.down.sql b/migrations/sqlite/000007_cdn_source_last_refreshed_at.down.sql new file mode 100644 index 0000000..33d0552 --- /dev/null +++ b/migrations/sqlite/000007_cdn_source_last_refreshed_at.down.sql @@ -0,0 +1,2 @@ +ALTER TABLE module_cdn_source +DROP COLUMN last_refreshed_at; diff --git a/migrations/sqlite/000007_cdn_source_last_refreshed_at.up.sql b/migrations/sqlite/000007_cdn_source_last_refreshed_at.up.sql new file mode 100644 index 0000000..d16bddd --- /dev/null +++ b/migrations/sqlite/000007_cdn_source_last_refreshed_at.up.sql @@ -0,0 +1,2 @@ +ALTER TABLE module_cdn_source +ADD COLUMN last_refreshed_at TEXT;