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;