Improve remote rule-set update
This commit is contained in:
parent
f7779ccea2
commit
c688b61c3b
6 changed files with 98 additions and 42 deletions
|
|
@ -46,7 +46,6 @@ type ConnectionRouterEx interface {
|
||||||
type RuleSet interface {
|
type RuleSet interface {
|
||||||
Name() string
|
Name() string
|
||||||
StartContext(ctx context.Context, startContext *HTTPStartContext) error
|
StartContext(ctx context.Context, startContext *HTTPStartContext) error
|
||||||
PostStart() error
|
|
||||||
Metadata() RuleSetMetadata
|
Metadata() RuleSetMetadata
|
||||||
ExtractIPSet() []*netipx.IPSet
|
ExtractIPSet() []*netipx.IPSet
|
||||||
IncRef()
|
IncRef()
|
||||||
|
|
|
||||||
|
|
@ -40,6 +40,7 @@ type Router struct {
|
||||||
leaseFiles []string
|
leaseFiles []string
|
||||||
ruleSets []adapter.RuleSet
|
ruleSets []adapter.RuleSet
|
||||||
ruleSetMap map[string]adapter.RuleSet
|
ruleSetMap map[string]adapter.RuleSet
|
||||||
|
ruleSetUpdater *R.RuleSetUpdater
|
||||||
processSearcher process.Searcher
|
processSearcher process.Searcher
|
||||||
processCache freelru.Cache[processCacheKey, processCacheEntry]
|
processCache freelru.Cache[processCacheKey, processCacheEntry]
|
||||||
neighborResolver adapter.NeighborResolver
|
neighborResolver adapter.NeighborResolver
|
||||||
|
|
@ -156,6 +157,7 @@ func (r *Router) Start(stage adapter.StartStage) error {
|
||||||
if startContext != nil {
|
if startContext != nil {
|
||||||
startContext.Close()
|
startContext.Close()
|
||||||
}
|
}
|
||||||
|
r.ruleSetUpdater = R.NewRuleSetUpdater(r.ctx, r.ruleSets)
|
||||||
r.network.Initialize(r.ruleSets)
|
r.network.Initialize(r.ruleSets)
|
||||||
needFindProcess := r.needFindProcess
|
needFindProcess := r.needFindProcess
|
||||||
for _, ruleSet := range r.ruleSets {
|
for _, ruleSet := range r.ruleSets {
|
||||||
|
|
@ -201,13 +203,8 @@ func (r *Router) Start(stage adapter.StartStage) error {
|
||||||
return E.Cause(err, "initialize rule[", i, "]")
|
return E.Cause(err, "initialize rule[", i, "]")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
for _, ruleSet := range r.ruleSets {
|
if r.ruleSetUpdater != nil {
|
||||||
monitor.Start("post start rule_set[", ruleSet.Name(), "]")
|
r.ruleSetUpdater.Start()
|
||||||
err := ruleSet.PostStart()
|
|
||||||
monitor.Finish()
|
|
||||||
if err != nil {
|
|
||||||
return E.Cause(err, "post start rule_set[", ruleSet.Name(), "]")
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
r.started = true
|
r.started = true
|
||||||
return nil
|
return nil
|
||||||
|
|
@ -237,6 +234,13 @@ func (r *Router) Close() error {
|
||||||
})
|
})
|
||||||
monitor.Finish()
|
monitor.Finish()
|
||||||
}
|
}
|
||||||
|
if r.ruleSetUpdater != nil {
|
||||||
|
monitor.Start("close rule-set updater")
|
||||||
|
err = E.Append(err, r.ruleSetUpdater.Close(), func(err error) error {
|
||||||
|
return E.Cause(err, "close rule-set updater")
|
||||||
|
})
|
||||||
|
monitor.Finish()
|
||||||
|
}
|
||||||
for i, ruleSet := range r.ruleSets {
|
for i, ruleSet := range r.ruleSets {
|
||||||
monitor.Start("close rule-set[", i, "]")
|
monitor.Start("close rule-set[", i, "]")
|
||||||
err = E.Append(err, ruleSet.Close(), func(err error) error {
|
err = E.Append(err, ruleSet.Close(), func(err error) error {
|
||||||
|
|
|
||||||
|
|
@ -24,10 +24,6 @@ func (f *fakeRuleSet) StartContext(context.Context, *adapter.HTTPStartContext) e
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (f *fakeRuleSet) PostStart() error {
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (f *fakeRuleSet) Metadata() adapter.RuleSetMetadata {
|
func (f *fakeRuleSet) Metadata() adapter.RuleSetMetadata {
|
||||||
return adapter.RuleSetMetadata{}
|
return adapter.RuleSetMetadata{}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -153,10 +153,6 @@ func (s *LocalRuleSet) reloadRules(headlessRules []option.HeadlessRule) error {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *LocalRuleSet) PostStart() error {
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (s *LocalRuleSet) Metadata() adapter.RuleSetMetadata {
|
func (s *LocalRuleSet) Metadata() adapter.RuleSetMetadata {
|
||||||
s.access.RLock()
|
s.access.RLock()
|
||||||
defer s.access.RUnlock()
|
defer s.access.RUnlock()
|
||||||
|
|
|
||||||
|
|
@ -5,7 +5,6 @@ import (
|
||||||
"context"
|
"context"
|
||||||
"io"
|
"io"
|
||||||
"net/http"
|
"net/http"
|
||||||
"runtime"
|
|
||||||
"strings"
|
"strings"
|
||||||
"sync"
|
"sync"
|
||||||
"sync/atomic"
|
"sync/atomic"
|
||||||
|
|
@ -43,7 +42,6 @@ type RemoteRuleSet struct {
|
||||||
metadata adapter.RuleSetMetadata
|
metadata adapter.RuleSetMetadata
|
||||||
lastUpdated time.Time
|
lastUpdated time.Time
|
||||||
lastEtag string
|
lastEtag string
|
||||||
updateTicker *time.Ticker
|
|
||||||
cacheFile adapter.CacheFile
|
cacheFile adapter.CacheFile
|
||||||
pauseManager pause.Manager
|
pauseManager pause.Manager
|
||||||
callbacks list.List[adapter.RuleSetUpdateCallback]
|
callbacks list.List[adapter.RuleSetUpdateCallback]
|
||||||
|
|
@ -102,12 +100,6 @@ func (s *RemoteRuleSet) StartContext(ctx context.Context, startContext *adapter.
|
||||||
return E.Cause(err, "initial rule-set: ", s.options.Tag)
|
return E.Cause(err, "initial rule-set: ", s.options.Tag)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
s.updateTicker = time.NewTicker(s.updateInterval)
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (s *RemoteRuleSet) PostStart() error {
|
|
||||||
go s.loopUpdate()
|
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -197,21 +189,6 @@ func (s *RemoteRuleSet) loadBytes(content []byte) error {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *RemoteRuleSet) loopUpdate() {
|
|
||||||
if time.Since(s.lastUpdated) > s.updateInterval {
|
|
||||||
s.updateOnce()
|
|
||||||
}
|
|
||||||
for {
|
|
||||||
runtime.GC()
|
|
||||||
select {
|
|
||||||
case <-s.ctx.Done():
|
|
||||||
return
|
|
||||||
case <-s.updateTicker.C:
|
|
||||||
s.updateOnce()
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (s *RemoteRuleSet) updateOnce() {
|
func (s *RemoteRuleSet) updateOnce() {
|
||||||
err := s.fetch(s.ctx, false)
|
err := s.fetch(s.ctx, false)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|
@ -312,9 +289,6 @@ func (s *RemoteRuleSet) resolveTransport() (adapter.HTTPTransport, error) {
|
||||||
func (s *RemoteRuleSet) Close() error {
|
func (s *RemoteRuleSet) Close() error {
|
||||||
s.rules = nil
|
s.rules = nil
|
||||||
s.cancel()
|
s.cancel()
|
||||||
if s.updateTicker != nil {
|
|
||||||
s.updateTicker.Stop()
|
|
||||||
}
|
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
87
route/rule/rule_set_updater.go
Normal file
87
route/rule/rule_set_updater.go
Normal file
|
|
@ -0,0 +1,87 @@
|
||||||
|
package rule
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"runtime"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/sagernet/sing-box/adapter"
|
||||||
|
)
|
||||||
|
|
||||||
|
type RuleSetUpdater struct {
|
||||||
|
ctx context.Context
|
||||||
|
cancel context.CancelFunc
|
||||||
|
ruleSets []*RemoteRuleSet
|
||||||
|
}
|
||||||
|
|
||||||
|
func NewRuleSetUpdater(ctx context.Context, ruleSets []adapter.RuleSet) *RuleSetUpdater {
|
||||||
|
var remoteRuleSets []*RemoteRuleSet
|
||||||
|
for _, ruleSet := range ruleSets {
|
||||||
|
remoteRuleSet, isRemote := ruleSet.(*RemoteRuleSet)
|
||||||
|
if isRemote {
|
||||||
|
remoteRuleSets = append(remoteRuleSets, remoteRuleSet)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if len(remoteRuleSets) == 0 {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
ctx, cancel := context.WithCancel(ctx)
|
||||||
|
return &RuleSetUpdater{
|
||||||
|
ctx: ctx,
|
||||||
|
cancel: cancel,
|
||||||
|
ruleSets: remoteRuleSets,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (u *RuleSetUpdater) Start() {
|
||||||
|
go u.loopUpdate()
|
||||||
|
}
|
||||||
|
|
||||||
|
func (u *RuleSetUpdater) Close() error {
|
||||||
|
u.cancel()
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (u *RuleSetUpdater) loopUpdate() {
|
||||||
|
nextUpdates := make([]time.Time, len(u.ruleSets))
|
||||||
|
for i, ruleSet := range u.ruleSets {
|
||||||
|
nextUpdates[i] = ruleSet.lastUpdated.Add(ruleSet.updateInterval)
|
||||||
|
}
|
||||||
|
timer := time.NewTimer(0)
|
||||||
|
defer timer.Stop()
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-u.ctx.Done():
|
||||||
|
return
|
||||||
|
case <-timer.C:
|
||||||
|
}
|
||||||
|
now := time.Now()
|
||||||
|
var updated bool
|
||||||
|
for i, ruleSet := range u.ruleSets {
|
||||||
|
if now.Before(nextUpdates[i]) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
ruleSet.updateOnce()
|
||||||
|
nextUpdates[i] = now.Add(ruleSet.updateInterval)
|
||||||
|
updated = true
|
||||||
|
}
|
||||||
|
if updated {
|
||||||
|
runtime.GC()
|
||||||
|
}
|
||||||
|
timer.Reset(waitUntilNext(nextUpdates))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func waitUntilNext(nextUpdates []time.Time) time.Duration {
|
||||||
|
next := nextUpdates[0]
|
||||||
|
for _, nextUpdate := range nextUpdates[1:] {
|
||||||
|
if nextUpdate.Before(next) {
|
||||||
|
next = nextUpdate
|
||||||
|
}
|
||||||
|
}
|
||||||
|
wait := time.Until(next)
|
||||||
|
if wait < 0 {
|
||||||
|
return 0
|
||||||
|
}
|
||||||
|
return wait
|
||||||
|
}
|
||||||
Loading…
Add table
Reference in a new issue