From 8425c428585e98cbac22fd29b14dfb4436d9a0c7 Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Thu, 7 May 2026 15:43:27 -0700 Subject: [PATCH] feat(s3/lifecycle): event router + schedule (Phase 3 PR-C) (#9355) feat(s3/lifecycle): event router + DueTime schedule Router consumes per-shard reader events, looks up matching ActionKeys via the engine's BucketActionKeys index, and emits Matches with DueTime = event_time + action.Delay. Evaluation runs at DueTime so the age gate passes for fresh events; the dispatcher's identity-CAS catches drift. Schedule is a min-heap by DueTime; duplicates allowed (RPC CAS handles the redundant dispatch as NOOP_RESOLVED). BucketActionKeys accessor added to engine.Snapshot. --- weed/s3api/s3lifecycle/engine/engine.go | 12 ++ weed/s3api/s3lifecycle/router/router.go | 165 ++++++++++++++++++ weed/s3api/s3lifecycle/router/router_test.go | 157 +++++++++++++++++ weed/s3api/s3lifecycle/router/schedule.go | 76 ++++++++ .../s3api/s3lifecycle/router/schedule_test.go | 86 +++++++++ 5 files changed, 496 insertions(+) create mode 100644 weed/s3api/s3lifecycle/router/router.go create mode 100644 weed/s3api/s3lifecycle/router/router_test.go create mode 100644 weed/s3api/s3lifecycle/router/schedule.go create mode 100644 weed/s3api/s3lifecycle/router/schedule_test.go diff --git a/weed/s3api/s3lifecycle/engine/engine.go b/weed/s3api/s3lifecycle/engine/engine.go index 7ac0b8856..0c6b40f04 100644 --- a/weed/s3api/s3lifecycle/engine/engine.go +++ b/weed/s3api/s3lifecycle/engine/engine.go @@ -118,3 +118,15 @@ func (s *Snapshot) BucketVersioned(bucket string) bool { } return false } + +// BucketActionKeys returns a defensive copy of the action keys for bucket, +// or nil if the bucket has no compiled actions in this snapshot. +func (s *Snapshot) BucketActionKeys(bucket string) []s3lifecycle.ActionKey { + bi, ok := s.buckets[bucket] + if !ok { + return nil + } + out := make([]s3lifecycle.ActionKey, len(bi.actionKeys)) + copy(out, bi.actionKeys) + return out +} diff --git a/weed/s3api/s3lifecycle/router/router.go b/weed/s3api/s3lifecycle/router/router.go new file mode 100644 index 000000000..227cb8cbf --- /dev/null +++ b/weed/s3api/s3lifecycle/router/router.go @@ -0,0 +1,165 @@ +package router + +import ( + "strings" + "time" + + "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" + "github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants" + "github.com/seaweedfs/seaweedfs/weed/s3api/s3lifecycle" + "github.com/seaweedfs/seaweedfs/weed/s3api/s3lifecycle/engine" + "github.com/seaweedfs/seaweedfs/weed/s3api/s3lifecycle/reader" +) + +// Match is one (event, action) pair where EvaluateAction fired. The +// dispatcher runs `LifecycleDelete` at DueTime; identity-CAS in the RPC +// guards against drift between schedule time and dispatch time. +type Match struct { + Key s3lifecycle.ActionKey + Action *engine.CompiledAction + Result s3lifecycle.EvalResult + EventTs time.Time + DueTime time.Time + Bucket string + ObjectKey string + VersionID string + Identity *EntryIdentity +} + +// EntryIdentity is the schedule-time CAS witness; the dispatcher serializes +// it into the LifecycleDelete request. The fields mirror +// s3_lifecycle_pb.EntryIdentity but stay in-package so the router doesn't +// pull a proto dependency. +type EntryIdentity struct { + MtimeNs int64 + Size int64 + HeadFid string +} + +// Route returns the matches that fire for ev against snap. Only EVENT_DRIVEN +// actions on this bucket are considered; actions in SCAN_AT_DATE or DISABLED +// modes are out-of-band of the event stream. Inactive actions +// (BootstrapComplete=false) are also skipped. +func Route(snap *engine.Snapshot, ev *reader.Event, now time.Time) []Match { + if snap == nil || ev == nil { + return nil + } + keys := snap.BucketActionKeys(ev.Bucket) + if len(keys) == 0 { + return nil + } + versioned := snap.BucketVersioned(ev.Bucket) + info := buildObjectInfo(ev, versioned) + if info == nil { + return nil + } + eventTime := time.Unix(0, ev.TsNs) + + var matches []Match + for _, key := range keys { + action := snap.Action(key) + if action == nil || !action.IsActive() { + continue + } + if action.Mode != engine.ModeEventDriven { + continue + } + // Evaluate at the scheduled dispatch time, not the event time: + // ExpirationDays gates on now >= modtime + N, which is exactly + // what holds when we actually dispatch. The dispatcher's + // identity-CAS catches drift if the object changes meanwhile. + dueTime := eventTime.Add(action.Delay) + res := s3lifecycle.EvaluateAction(action.Rule, key.ActionKind, info, dueTime) + if res.Action == s3lifecycle.ActionNone { + continue + } + matches = append(matches, Match{ + Key: key, + Action: action, + Result: res, + EventTs: eventTime, + DueTime: dueTime, + Bucket: ev.Bucket, + ObjectKey: ev.Key, + Identity: buildIdentity(ev), + }) + } + return matches +} + +// buildObjectInfo derives a non-versioned ObjectInfo from a meta-log event. +// Versioned-bucket semantics (IsLatest, NumVersions, NoncurrentIndex, +// IsDeleteMarker for noncurrent versions) require listing siblings and land +// in Phase 5; for now an event on a versioned bucket is treated as +// IsLatest=true with the same caveat that the LifecycleDelete RPC's +// identity-CAS catches stale schedules. +func buildObjectInfo(ev *reader.Event, versioned bool) *s3lifecycle.ObjectInfo { + entry := ev.NewEntry + if entry == nil { + entry = ev.OldEntry + } + if entry == nil { + return nil + } + info := &s3lifecycle.ObjectInfo{ + Key: ev.Key, + IsLatest: true, + NumVersions: 1, + } + if entry.Attributes != nil { + info.ModTime = time.Unix(entry.Attributes.Mtime, 0) + info.Size = int64(entry.Attributes.FileSize) + } + if tags := extractTags(entry.Extended); len(tags) > 0 { + info.Tags = tags + } + if isDeleteMarkerEntry(entry) { + info.IsDeleteMarker = true + } + _ = versioned + return info +} + +// buildIdentity captures the entry's schedule-time fingerprint for the CAS +// witness. Returns nil if the event has no entry to fingerprint (deletes). +func buildIdentity(ev *reader.Event) *EntryIdentity { + entry := ev.NewEntry + if entry == nil { + return nil + } + id := &EntryIdentity{} + if entry.Attributes != nil { + id.MtimeNs = entry.Attributes.Mtime + id.Size = int64(entry.Attributes.FileSize) + } + if len(entry.GetChunks()) > 0 { + id.HeadFid = entry.GetChunks()[0].FileId + } + return id +} + +func extractTags(ext map[string][]byte) map[string]string { + if len(ext) == 0 { + return nil + } + prefix := s3_constants.AmzObjectTagging + "-" + var out map[string]string + for k, v := range ext { + if !strings.HasPrefix(k, prefix) { + continue + } + if out == nil { + out = map[string]string{} + } + out[k[len(prefix):]] = string(v) + } + return out +} + +func isDeleteMarkerEntry(entry *filer_pb.Entry) bool { + if entry == nil || len(entry.Extended) == 0 { + return false + } + v, ok := entry.Extended[s3_constants.ExtDeleteMarkerKey] + return ok && len(v) == 1 && v[0] == 1 +} diff --git a/weed/s3api/s3lifecycle/router/router_test.go b/weed/s3api/s3lifecycle/router/router_test.go new file mode 100644 index 000000000..11db7cba5 --- /dev/null +++ b/weed/s3api/s3lifecycle/router/router_test.go @@ -0,0 +1,157 @@ +package router + +import ( + "testing" + "time" + + "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" + "github.com/seaweedfs/seaweedfs/weed/s3api/s3lifecycle" + "github.com/seaweedfs/seaweedfs/weed/s3api/s3lifecycle/engine" + "github.com/seaweedfs/seaweedfs/weed/s3api/s3lifecycle/reader" +) + +func compileWith(rule *s3lifecycle.Rule, prior map[s3lifecycle.ActionKey]engine.PriorState) *engine.Snapshot { + e := engine.New() + return e.Compile([]engine.CompileInput{{Bucket: "bk", Rules: []*s3lifecycle.Rule{rule}}}, + engine.CompileOptions{PriorStates: prior}) +} + +func activatedPrior(rule *s3lifecycle.Rule) map[s3lifecycle.ActionKey]engine.PriorState { + prior := map[s3lifecycle.ActionKey]engine.PriorState{} + hash := s3lifecycle.RuleHash(rule) + for _, k := range s3lifecycle.RuleActionKinds(rule) { + prior[s3lifecycle.ActionKey{Bucket: "bk", RuleHash: hash, ActionKind: k}] = engine.PriorState{ + BootstrapComplete: true, + Mode: engine.ModeEventDriven, + } + } + return prior +} + +func eventCreate(bucket, key string, modTimeS, size int64, ts int64) *reader.Event { + return &reader.Event{ + TsNs: ts, + Bucket: bucket, + Key: key, + NewEntry: &filer_pb.Entry{ + Name: key, + Attributes: &filer_pb.FuseAttributes{ + Mtime: modTimeS, + FileSize: uint64(size), + }, + }, + } +} + +func TestRouteNoSnapshotNoMatches(t *testing.T) { + if got := Route(nil, eventCreate("bk", "k", 0, 1, 1), time.Now()); got != nil { + t.Fatalf("nil snap should yield nil, got %v", got) + } +} + +func TestRouteMissingBucketNoMatches(t *testing.T) { + rule := &s3lifecycle.Rule{ID: "r", Status: s3lifecycle.StatusEnabled, ExpirationDays: 1} + snap := compileWith(rule, activatedPrior(rule)) + ev := eventCreate("other-bucket", "k", 0, 1, 1) + if got := Route(snap, ev, time.Now()); got != nil { + t.Fatalf("foreign bucket should yield nil, got %v", got) + } +} + +func TestRouteInactiveSkipped(t *testing.T) { + rule := &s3lifecycle.Rule{ID: "r", Status: s3lifecycle.StatusEnabled, ExpirationDays: 1} + // PriorStates omitted => BootstrapComplete=false => action stays inactive. + snap := compileWith(rule, nil) + + now := time.Now() + old := now.Add(-48 * time.Hour) // past 1-day expiration + ev := eventCreate("bk", "k", old.Unix(), 1, old.UnixNano()) + + if got := Route(snap, ev, now); got != nil { + t.Fatalf("inactive action should not match, got %v", got) + } +} + +func TestRouteExpirationDaysFires(t *testing.T) { + rule := &s3lifecycle.Rule{ID: "r", Status: s3lifecycle.StatusEnabled, ExpirationDays: 1} + snap := compileWith(rule, activatedPrior(rule)) + + now := time.Now() + old := now.Add(-48 * time.Hour) // past 1-day expiration + ev := eventCreate("bk", "k", old.Unix(), 1, old.UnixNano()) + + matches := Route(snap, ev, now) + if len(matches) != 1 { + t.Fatalf("expected 1 match, got %+v", matches) + } + m := matches[0] + if m.Bucket != "bk" || m.ObjectKey != "k" { + t.Fatalf("unexpected match shape: %+v", m) + } + if m.Result.Action != s3lifecycle.ActionDeleteObject { + t.Fatalf("action=%v, want DeleteObject", m.Result.Action) + } + if m.DueTime.Before(m.EventTs) { + t.Fatalf("DueTime < EventTs: %v < %v", m.DueTime, m.EventTs) + } +} + +func TestRouteFreshObjectSchedulesInFuture(t *testing.T) { + // A fresh object DOES route — the action is scheduled for ModTime+N days. + // EvaluateAction is called at the scheduled dispatch time so the age + // gate passes; the dispatcher waits until DueTime to actually run. + rule := &s3lifecycle.Rule{ID: "r", Status: s3lifecycle.StatusEnabled, ExpirationDays: 7} + snap := compileWith(rule, activatedPrior(rule)) + + now := time.Now() + ev := eventCreate("bk", "k", now.Unix(), 1, now.UnixNano()) + + matches := Route(snap, ev, now) + if len(matches) != 1 { + t.Fatalf("expected 1 match (scheduled), got %v", matches) + } + if !matches[0].DueTime.After(now.Add(6 * 24 * time.Hour)) { + t.Fatalf("DueTime=%v, want ~7 days from now", matches[0].DueTime) + } +} + +func TestRouteRespectsPrefixFilter(t *testing.T) { + rule := &s3lifecycle.Rule{ + ID: "r", Status: s3lifecycle.StatusEnabled, + Prefix: "logs/", ExpirationDays: 1, + } + snap := compileWith(rule, activatedPrior(rule)) + now := time.Now() + old := now.Add(-48 * time.Hour) + + // Out of prefix: no match. + ev := eventCreate("bk", "data/file", old.Unix(), 1, old.UnixNano()) + if got := Route(snap, ev, now); got != nil { + t.Fatalf("out-of-prefix should not match, got %v", got) + } + + // In prefix: matches. + ev = eventCreate("bk", "logs/file", old.Unix(), 1, old.UnixNano()) + if got := Route(snap, ev, now); len(got) != 1 { + t.Fatalf("in-prefix should match, got %v", got) + } +} + +func TestRouteIdentityCapturedForNewEntry(t *testing.T) { + rule := &s3lifecycle.Rule{ID: "r", Status: s3lifecycle.StatusEnabled, ExpirationDays: 1} + snap := compileWith(rule, activatedPrior(rule)) + + now := time.Now() + old := now.Add(-48 * time.Hour) + ev := eventCreate("bk", "k", old.Unix(), 42, old.UnixNano()) + ev.NewEntry.Chunks = []*filer_pb.FileChunk{{FileId: "1,abc"}} + + matches := Route(snap, ev, now) + if len(matches) != 1 { + t.Fatalf("expected 1 match, got %v", matches) + } + id := matches[0].Identity + if id == nil || id.Size != 42 || id.HeadFid != "1,abc" { + t.Fatalf("identity capture: %+v", id) + } +} diff --git a/weed/s3api/s3lifecycle/router/schedule.go b/weed/s3api/s3lifecycle/router/schedule.go new file mode 100644 index 000000000..50e97b46c --- /dev/null +++ b/weed/s3api/s3lifecycle/router/schedule.go @@ -0,0 +1,76 @@ +package router + +import ( + "container/heap" + "sync" + "time" +) + +// Schedule is the worker's in-memory pending list, ordered by DueTime. +// The dispatcher polls Drain on each tick to retrieve due Matches. +// +// Duplicates are allowed: a re-schedule for the same (ActionKey, ObjectKey) +// before the prior dispatch ran results in two heap entries. The +// LifecycleDelete RPC's identity-CAS makes the second dispatch a no-op +// (NOOP_RESOLVED with STALE_IDENTITY), so dedup at insert time would only +// be a micro-optimization at the cost of an extra map. +type Schedule struct { + mu sync.Mutex + h scheduleHeap +} + +func NewSchedule() *Schedule { + return &Schedule{} +} + +// Add enqueues a Match. +func (s *Schedule) Add(m Match) { + s.mu.Lock() + defer s.mu.Unlock() + heap.Push(&s.h, m) +} + +// Len returns the number of pending Matches. +func (s *Schedule) Len() int { + s.mu.Lock() + defer s.mu.Unlock() + return s.h.Len() +} + +// NextDue returns the DueTime of the earliest pending Match. ok=false if +// the schedule is empty. +func (s *Schedule) NextDue() (time.Time, bool) { + s.mu.Lock() + defer s.mu.Unlock() + if s.h.Len() == 0 { + return time.Time{}, false + } + return s.h[0].DueTime, true +} + +// Drain pops and returns all Matches whose DueTime <= now, in DueTime order. +// Subsequent calls return only newly-due Matches. +func (s *Schedule) Drain(now time.Time) []Match { + s.mu.Lock() + defer s.mu.Unlock() + var out []Match + for s.h.Len() > 0 && !s.h[0].DueTime.After(now) { + out = append(out, heap.Pop(&s.h).(Match)) + } + return out +} + +// scheduleHeap implements heap.Interface ordered by Match.DueTime ascending. +type scheduleHeap []Match + +func (h scheduleHeap) Len() int { return len(h) } +func (h scheduleHeap) Less(i, j int) bool { return h[i].DueTime.Before(h[j].DueTime) } +func (h scheduleHeap) Swap(i, j int) { h[i], h[j] = h[j], h[i] } +func (h *scheduleHeap) Push(x interface{}) { *h = append(*h, x.(Match)) } +func (h *scheduleHeap) Pop() interface{} { + old := *h + n := len(old) + x := old[n-1] + *h = old[:n-1] + return x +} diff --git a/weed/s3api/s3lifecycle/router/schedule_test.go b/weed/s3api/s3lifecycle/router/schedule_test.go new file mode 100644 index 000000000..f6cd0827b --- /dev/null +++ b/weed/s3api/s3lifecycle/router/schedule_test.go @@ -0,0 +1,86 @@ +package router + +import ( + "testing" + "time" + + "github.com/seaweedfs/seaweedfs/weed/s3api/s3lifecycle" +) + +func mkMatch(due time.Time, key string) Match { + return Match{ + Key: s3lifecycle.ActionKey{Bucket: key}, + DueTime: due, + ObjectKey: key, + } +} + +func TestScheduleEmpty(t *testing.T) { + s := NewSchedule() + if s.Len() != 0 { + t.Fatal("Len != 0") + } + if _, ok := s.NextDue(); ok { + t.Fatal("NextDue ok=true on empty") + } + if got := s.Drain(time.Now()); got != nil { + t.Fatalf("Drain on empty returned %v", got) + } +} + +func TestScheduleOrderedByDueTime(t *testing.T) { + s := NewSchedule() + t0 := time.Now() + s.Add(mkMatch(t0.Add(3*time.Second), "c")) + s.Add(mkMatch(t0.Add(1*time.Second), "a")) + s.Add(mkMatch(t0.Add(2*time.Second), "b")) + + if s.Len() != 3 { + t.Fatalf("Len=%d, want 3", s.Len()) + } + due, ok := s.NextDue() + if !ok || !due.Equal(t0.Add(1*time.Second)) { + t.Fatalf("NextDue=%v ok=%v", due, ok) + } + + got := s.Drain(t0.Add(2 * time.Second)) + if len(got) != 2 || got[0].ObjectKey != "a" || got[1].ObjectKey != "b" { + t.Fatalf("Drain order: %+v", got) + } + if s.Len() != 1 { + t.Fatalf("Len after drain=%d, want 1", s.Len()) + } + + got = s.Drain(t0.Add(5 * time.Second)) + if len(got) != 1 || got[0].ObjectKey != "c" { + t.Fatalf("Drain rest: %+v", got) + } + if s.Len() != 0 { + t.Fatal("Len != 0 after final drain") + } +} + +func TestScheduleDrainBoundaryInclusive(t *testing.T) { + // DueTime exactly equal to now is drainable (<=). + s := NewSchedule() + t0 := time.Now() + s.Add(mkMatch(t0, "a")) + got := s.Drain(t0) + if len(got) != 1 { + t.Fatalf("expected boundary-inclusive drain, got %d", len(got)) + } +} + +func TestScheduleAllowsDuplicates(t *testing.T) { + s := NewSchedule() + t0 := time.Now() + s.Add(mkMatch(t0, "a")) + s.Add(mkMatch(t0, "a")) + if s.Len() != 2 { + t.Fatalf("dup count=%d, want 2", s.Len()) + } + got := s.Drain(t0) + if len(got) != 2 { + t.Fatalf("Drain dup count=%d, want 2", len(got)) + } +}