mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-01 14:46:36 +08:00
feat(core): add config resolution barrier and plugin gating to App
App 新增 WithConfigSource / WithConfigDecl / Prepare / ShutdownTimeout / SetShutdownTimeout;Use 收集门禁插件的提前声明,调和循环内求值门禁并跳过 被关闭的插件,使组合根无需再跨插件读配置选实现。未注入配置源的 App 保持 原行为,被门禁但无配置源则 fail fast 点名原因。
This commit is contained in:
+154
-13
@@ -68,22 +68,54 @@ func WithShutdownTimeout(timeout time.Duration) AppOption {
|
||||
}
|
||||
}
|
||||
|
||||
// WithConfigSource installs the raw configuration source adapter, typically built by an
|
||||
// infrastructure package outside the kernel, before any plugin is applied.
|
||||
func WithConfigSource(src ConfigSource) AppOption {
|
||||
return func(a *App) {
|
||||
if src == nil {
|
||||
return
|
||||
}
|
||||
// Installed during Prepare so the option order, including WithContext, is irrelevant.
|
||||
a.configSource = src
|
||||
}
|
||||
}
|
||||
|
||||
// WithConfigDecl lets the composition root declare the configuration it reads itself,
|
||||
// so host-level values take part in conflict validation and the redacted report. The
|
||||
// bindings are registered during Prepare, so option order does not matter.
|
||||
func WithConfigDecl(pluginID string, bindings ...ConfigBinding) AppOption {
|
||||
return func(a *App) {
|
||||
if len(bindings) == 0 {
|
||||
return
|
||||
}
|
||||
if a.hostDeclOwner == "" {
|
||||
a.hostDeclOwner = pluginID
|
||||
}
|
||||
a.hostDeclBindings = append(a.hostDeclBindings, bindings...)
|
||||
}
|
||||
}
|
||||
|
||||
// App is the unified assembly entrypoint and runtime aspect dispatcher of the Cordis micro-kernel.
|
||||
// It manages plugin collection, dependency mounting, migration execution, profile-based driver startup,
|
||||
// and graceful signal-driven LIFO shutdown.
|
||||
type App struct {
|
||||
mu sync.RWMutex
|
||||
ctx *Context
|
||||
profile Profile
|
||||
plugins []Plugin
|
||||
pluginMap map[string]Plugin
|
||||
fibers []*Fiber
|
||||
fiberMap map[string]*Fiber
|
||||
applied bool
|
||||
running bool
|
||||
startedDrivers []Driver
|
||||
migrationEngine MigrationEngine
|
||||
shutdownTimeout time.Duration
|
||||
mu sync.RWMutex
|
||||
ctx *Context
|
||||
profile Profile
|
||||
plugins []Plugin
|
||||
pluginMap map[string]Plugin
|
||||
fibers []*Fiber
|
||||
fiberMap map[string]*Fiber
|
||||
applied bool
|
||||
running bool
|
||||
startedDrivers []Driver
|
||||
migrationEngine MigrationEngine
|
||||
shutdownTimeout time.Duration
|
||||
configSource ConfigSource
|
||||
hostDeclOwner string
|
||||
hostDeclBindings []ConfigBinding
|
||||
prepared bool
|
||||
applyErr error
|
||||
}
|
||||
|
||||
// NewApp creates a new Cordis application instance with default options.
|
||||
@@ -162,6 +194,11 @@ func (a *App) Use(plugins ...Plugin) *App {
|
||||
a.fiberMap[name] = f
|
||||
}
|
||||
a.pluginMap[name] = p
|
||||
|
||||
if gated, ok := p.(ConfigGatedPlugin); ok && a.applyErr == nil {
|
||||
// Gates are evaluated before Apply, so their keys must be declared at mount time.
|
||||
a.applyErr = a.ctx.Config().Declare(name, gated.DeclareConfig()...)
|
||||
}
|
||||
}
|
||||
|
||||
return a
|
||||
@@ -227,10 +264,29 @@ func (a *App) Reconcile() error {
|
||||
}
|
||||
|
||||
func (a *App) reconcileLocked() error {
|
||||
if err := a.prepareLocked(); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
for {
|
||||
progress := false
|
||||
for _, f := range a.fibers {
|
||||
if f.State() == FiberPending && f.DependenciesSatisfied(a.ctx) {
|
||||
if f.State() != FiberPending {
|
||||
continue
|
||||
}
|
||||
|
||||
gated, skip, err := a.evaluateGateLocked(f)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if gated && skip {
|
||||
if err := f.Skip(); err != nil {
|
||||
return fmt.Errorf("core: skip gated fiber %q failed: %w", f.Name(), err)
|
||||
}
|
||||
continue
|
||||
}
|
||||
|
||||
if f.DependenciesSatisfied(a.ctx) {
|
||||
if err := f.Load(); err != nil {
|
||||
return fmt.Errorf("core: load fiber %q failed: %w", f.Name(), err)
|
||||
}
|
||||
@@ -255,6 +311,24 @@ func (a *App) reconcileLocked() error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// evaluateGateLocked reports whether a configuration-gated plugin is excluded by the
|
||||
// resolved values. Plugins that do not implement the gate interface are never skipped.
|
||||
func (a *App) evaluateGateLocked(f *Fiber) (gated bool, skip bool, err error) {
|
||||
gatedPlugin, ok := f.plugin.(ConfigGatedPlugin)
|
||||
if !ok {
|
||||
return false, false, nil
|
||||
}
|
||||
|
||||
view := a.ctx.Config()
|
||||
if !view.Resolved() {
|
||||
return true, false, fmt.Errorf(
|
||||
"core: plugin %q is configuration-gated but the App has no ConfigSource; "+
|
||||
"pass core.WithConfigSource or remove DeclareConfig", f.Name())
|
||||
}
|
||||
|
||||
return true, !gatedPlugin.ConfigEnabled(view), nil
|
||||
}
|
||||
|
||||
// ApplyPlugins applies all registered plugins on the application Context via reactive reconciliation.
|
||||
// It is idempotent and only applies plugins once per App instance.
|
||||
func (a *App) ApplyPlugins() error {
|
||||
@@ -264,11 +338,78 @@ func (a *App) ApplyPlugins() error {
|
||||
return nil
|
||||
}
|
||||
a.applied = true
|
||||
|
||||
declaredErr, prepareErr := a.applyErr, a.prepareLocked()
|
||||
a.mu.Unlock()
|
||||
|
||||
if declaredErr != nil {
|
||||
return declaredErr
|
||||
}
|
||||
if prepareErr != nil {
|
||||
return prepareErr
|
||||
}
|
||||
|
||||
return a.Reconcile()
|
||||
}
|
||||
|
||||
// Prepare resolves declared configuration and establishes the resolution barrier that
|
||||
// gates and plugin Bind calls depend on. It is idempotent and runs implicitly from
|
||||
// ApplyPlugins; callers that need resolved values earlier — for example to size a
|
||||
// shutdown budget — invoke it explicitly right after mounting plugins.
|
||||
func (a *App) Prepare() error {
|
||||
a.mu.Lock()
|
||||
defer a.mu.Unlock()
|
||||
|
||||
if a.applyErr != nil {
|
||||
return a.applyErr
|
||||
}
|
||||
return a.prepareLocked()
|
||||
}
|
||||
|
||||
// prepareLocked installs the injected source, registers host declarations and resolves
|
||||
// every declared key once. An App without a ConfigSource leaves configuration unused,
|
||||
// so kernel-level usage stays opt-in for embedders that configure nothing.
|
||||
func (a *App) prepareLocked() error {
|
||||
if a.prepared {
|
||||
return nil
|
||||
}
|
||||
if a.configSource == nil {
|
||||
a.prepared = true
|
||||
return nil
|
||||
}
|
||||
|
||||
config := a.ctx.Config()
|
||||
config.SetSource(a.configSource)
|
||||
|
||||
if err := config.Declare(a.hostDeclOwner, a.hostDeclBindings...); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := config.Resolve(); err != nil {
|
||||
return err
|
||||
}
|
||||
a.prepared = true
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// ShutdownTimeout returns the graceful shutdown budget for the application.
|
||||
func (a *App) ShutdownTimeout() time.Duration {
|
||||
a.mu.RLock()
|
||||
defer a.mu.RUnlock()
|
||||
return a.shutdownTimeout
|
||||
}
|
||||
|
||||
// SetShutdownTimeout replaces the graceful shutdown budget, ignoring non-positive
|
||||
// values so a missing configuration key can never shrink the kernel fallback to zero.
|
||||
func (a *App) SetShutdownTimeout(timeout time.Duration) *App {
|
||||
a.mu.Lock()
|
||||
defer a.mu.Unlock()
|
||||
if timeout > 0 {
|
||||
a.shutdownTimeout = timeout
|
||||
}
|
||||
return a
|
||||
}
|
||||
|
||||
// RunMigrations dispatches migration execution across all registered plugin migration entries.
|
||||
func (a *App) RunMigrations() error {
|
||||
entries := a.ctx.Migrations().Entries()
|
||||
|
||||
@@ -504,3 +504,93 @@ func TestAppIdempotencyAndErrorStates(t *testing.T) {
|
||||
assert.Contains(t, err.Error(), "sql migrate error")
|
||||
assert.False(t, app3.IsRunning())
|
||||
}
|
||||
|
||||
// newGateSource builds a configuration source whose only key decides the test gates.
|
||||
func newGateSource(enabled bool) *mapSource {
|
||||
return &mapSource{
|
||||
values: map[string]any{"gate.enabled": enabled},
|
||||
env: map[string]string{},
|
||||
}
|
||||
}
|
||||
|
||||
func TestAppPrepareResolvesThenGatesDuringReconcile(t *testing.T) {
|
||||
primary := &gatedPlugin{name: "cache", enabled: true}
|
||||
fallback := &gatedPlugin{name: "cache_memory", enabled: false}
|
||||
|
||||
app := core.NewApp(core.WithConfigSource(newGateSource(true)))
|
||||
app.Use(primary, fallback)
|
||||
require.NoError(t, app.Prepare())
|
||||
|
||||
cacheFiber, ok := app.Fiber("cache")
|
||||
require.True(t, ok)
|
||||
require.Equal(t, core.FiberPending, cacheFiber.State(), "Prepare only builds the resolution barrier")
|
||||
assert.True(t, app.Context().Config().Resolved())
|
||||
|
||||
require.NoError(t, app.Reconcile())
|
||||
|
||||
assert.Equal(t, core.FiberActive, cacheFiber.State())
|
||||
|
||||
memoryFiber, ok := app.Fiber("cache_memory")
|
||||
require.True(t, ok)
|
||||
assert.Equal(t, core.FiberSkipped, memoryFiber.State())
|
||||
assert.False(t, fallback.applied, "the gated-out provider must never reach Apply")
|
||||
}
|
||||
|
||||
func TestAppGatesPluginsMountedAfterPrepare(t *testing.T) {
|
||||
app := core.NewApp(core.WithConfigSource(newGateSource(true)))
|
||||
require.NoError(t, app.Prepare())
|
||||
|
||||
late := &gatedPlugin{name: "cache_memory", enabled: false}
|
||||
app.Use(late)
|
||||
require.NoError(t, app.Reconcile())
|
||||
|
||||
fiber, ok := app.Fiber("cache_memory")
|
||||
require.True(t, ok)
|
||||
assert.Equal(t, core.FiberSkipped, fiber.State(),
|
||||
"plugins mounted after Prepare must still be gated")
|
||||
}
|
||||
|
||||
func TestAppApplyPluginsGatesImplicitly(t *testing.T) {
|
||||
app := core.NewApp(core.WithConfigSource(newGateSource(false)))
|
||||
app.Use(&gatedPlugin{name: "cache", enabled: true})
|
||||
|
||||
require.NoError(t, app.ApplyPlugins())
|
||||
|
||||
fiber, ok := app.Fiber("cache")
|
||||
require.True(t, ok)
|
||||
assert.Equal(t, core.FiberSkipped, fiber.State(),
|
||||
"ApplyPlugins must resolve and gate without an explicit Prepare call")
|
||||
}
|
||||
|
||||
func TestAppPrepareReportsConfigurationErrors(t *testing.T) {
|
||||
src := &mapSource{
|
||||
values: map[string]any{"gate.enabled": "yes"},
|
||||
env: map[string]string{},
|
||||
}
|
||||
app := core.NewApp(core.WithConfigSource(src))
|
||||
app.Use(&gatedPlugin{name: "cache", enabled: true})
|
||||
|
||||
err := app.Prepare()
|
||||
require.Error(t, err)
|
||||
assert.Contains(t, err.Error(), "gate.enabled")
|
||||
}
|
||||
|
||||
func TestAppGatedPluginWithoutConfigSourceFailsFast(t *testing.T) {
|
||||
app := core.NewApp()
|
||||
app.Use(&gatedPlugin{name: "cache", enabled: true})
|
||||
|
||||
err := app.ApplyPlugins()
|
||||
require.Error(t, err)
|
||||
assert.Contains(t, err.Error(), "cache")
|
||||
assert.Contains(t, err.Error(), "ConfigSource")
|
||||
}
|
||||
|
||||
func TestAppSetShutdownTimeoutIgnoresNonPositive(t *testing.T) {
|
||||
app := core.NewApp()
|
||||
|
||||
app.SetShutdownTimeout(0)
|
||||
assert.Equal(t, 10*time.Second, app.ShutdownTimeout(), "zero must not shrink the kernel fallback")
|
||||
|
||||
app.SetShutdownTimeout(45 * time.Second)
|
||||
assert.Equal(t, 45*time.Second, app.ShutdownTimeout())
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user