mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-09-29 14:06:36 +08:00
feat(core): implement plugin fiber state machine and reactive dependency reconciler
This commit is contained in:
+80
-10
@@ -77,6 +77,8 @@ type App struct {
|
||||
profile Profile
|
||||
plugins []Plugin
|
||||
pluginMap map[string]Plugin
|
||||
fibers []*Fiber
|
||||
fiberMap map[string]*Fiber
|
||||
applied bool
|
||||
running bool
|
||||
startedDrivers []Driver
|
||||
@@ -90,6 +92,7 @@ func NewApp(opts ...AppOption) *App {
|
||||
ctx: NewContext(context.Background()),
|
||||
profile: ProfileAll,
|
||||
pluginMap: make(map[string]Plugin),
|
||||
fiberMap: make(map[string]*Fiber),
|
||||
shutdownTimeout: defaultShutdownTimeout,
|
||||
}
|
||||
|
||||
@@ -149,8 +152,14 @@ func (a *App) Use(plugins ...Plugin) *App {
|
||||
break
|
||||
}
|
||||
}
|
||||
if existingFiber, ok := a.fiberMap[name]; ok {
|
||||
existingFiber.plugin = p
|
||||
}
|
||||
} else {
|
||||
a.plugins = append(a.plugins, p)
|
||||
f := NewFiber(a.ctx, p)
|
||||
a.fibers = append(a.fibers, f)
|
||||
a.fiberMap[name] = f
|
||||
}
|
||||
a.pluginMap[name] = p
|
||||
}
|
||||
@@ -177,6 +186,25 @@ func (a *App) Plugin(name string) (Plugin, bool) {
|
||||
return p, ok
|
||||
}
|
||||
|
||||
// Fibers returns a copy of all plugin Fibers.
|
||||
func (a *App) Fibers() []*Fiber {
|
||||
a.mu.RLock()
|
||||
defer a.mu.RUnlock()
|
||||
|
||||
res := make([]*Fiber, len(a.fibers))
|
||||
copy(res, a.fibers)
|
||||
return res
|
||||
}
|
||||
|
||||
// Fiber retrieves a Fiber by its unique plugin name.
|
||||
func (a *App) Fiber(name string) (*Fiber, bool) {
|
||||
a.mu.RLock()
|
||||
defer a.mu.RUnlock()
|
||||
|
||||
f, ok := a.fiberMap[name]
|
||||
return f, ok
|
||||
}
|
||||
|
||||
// SetMigrationEngine sets the migration engine for the application.
|
||||
func (a *App) SetMigrationEngine(engine MigrationEngine) *App {
|
||||
a.mu.Lock()
|
||||
@@ -190,7 +218,44 @@ func (a *App) SetMigrationRunner(runner MigrationRunner) *App {
|
||||
return a.SetMigrationEngine(runner)
|
||||
}
|
||||
|
||||
// ApplyPlugins applies all registered plugins on the application Context.
|
||||
// Reconcile evaluates all pending Fibers and reactively transitions them to ACTIVE
|
||||
// as their declared dependencies become satisfied.
|
||||
func (a *App) Reconcile() error {
|
||||
a.mu.Lock()
|
||||
defer a.mu.Unlock()
|
||||
return a.reconcileLocked()
|
||||
}
|
||||
|
||||
func (a *App) reconcileLocked() error {
|
||||
for {
|
||||
progress := false
|
||||
for _, f := range a.fibers {
|
||||
if f.State() == FiberPending && f.DependenciesSatisfied(a.ctx) {
|
||||
if err := f.Load(); err != nil {
|
||||
return fmt.Errorf("core: load fiber %q failed: %w", f.Name(), err)
|
||||
}
|
||||
progress = true
|
||||
}
|
||||
}
|
||||
if !progress {
|
||||
break
|
||||
}
|
||||
}
|
||||
|
||||
var unsatisfied []string
|
||||
for _, f := range a.fibers {
|
||||
if f.State() == FiberPending {
|
||||
unsatisfied = append(unsatisfied, fmt.Sprintf("%s (waiting for %v)", f.Name(), f.Dependencies()))
|
||||
}
|
||||
}
|
||||
if len(unsatisfied) > 0 {
|
||||
return fmt.Errorf("core: unsatisfied dependencies for plugins: %s", strings.Join(unsatisfied, ", "))
|
||||
}
|
||||
|
||||
return 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 {
|
||||
a.mu.Lock()
|
||||
@@ -199,16 +264,9 @@ func (a *App) ApplyPlugins() error {
|
||||
return nil
|
||||
}
|
||||
a.applied = true
|
||||
plugins := make([]Plugin, len(a.plugins))
|
||||
copy(plugins, a.plugins)
|
||||
a.mu.Unlock()
|
||||
|
||||
for _, p := range plugins {
|
||||
if err := p.Apply(a.ctx); err != nil {
|
||||
return fmt.Errorf("core: apply plugin %q failed: %w", p.Name(), err)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
return a.Reconcile()
|
||||
}
|
||||
|
||||
// RunMigrations dispatches migration execution across all registered plugin migration entries.
|
||||
@@ -363,7 +421,19 @@ func (a *App) Stop(ctx ...context.Context) error {
|
||||
}
|
||||
}
|
||||
|
||||
// 2. Dispose context
|
||||
// 2. Unload fibers in reverse order
|
||||
a.mu.RLock()
|
||||
fibers := make([]*Fiber, len(a.fibers))
|
||||
copy(fibers, a.fibers)
|
||||
a.mu.RUnlock()
|
||||
|
||||
for i := len(fibers) - 1; i >= 0; i-- {
|
||||
if err := fibers[i].Unload(); err != nil {
|
||||
errs = append(errs, fmt.Errorf("core: unload fiber %s failed: %w", fibers[i].Name(), err))
|
||||
}
|
||||
}
|
||||
|
||||
// 3. Dispose root context
|
||||
if a.ctx != nil && !a.ctx.IsDisposed() {
|
||||
if err := a.ctx.Dispose(); err != nil {
|
||||
errs = append(errs, fmt.Errorf("core: dispose context failed: %w", err))
|
||||
|
||||
@@ -46,7 +46,7 @@ func (c *Container) remove(targetType reflect.Type) {
|
||||
delete(c.services, targetType)
|
||||
}
|
||||
|
||||
// Provide registers a typed service implementation into the Context's IoC container.
|
||||
// Provide registers a typed service implementation into the Context hierarchy's root IoC container.
|
||||
func Provide[T any](ctx *Context, service T) {
|
||||
if ctx == nil {
|
||||
panic("core: nil context provided to Provide")
|
||||
@@ -55,6 +55,25 @@ func Provide[T any](ctx *Context, service T) {
|
||||
panic("core: cannot provide nil service")
|
||||
}
|
||||
|
||||
targetType := reflect.TypeFor[T]()
|
||||
targetContainer := ctx.Root().Container()
|
||||
targetContainer.provide(targetType, service)
|
||||
|
||||
ctx.OnDispose(func() error {
|
||||
targetContainer.remove(targetType)
|
||||
return nil
|
||||
})
|
||||
}
|
||||
|
||||
// ProvideScoped registers a typed service implementation strictly in the local Context container.
|
||||
func ProvideScoped[T any](ctx *Context, service T) {
|
||||
if ctx == nil {
|
||||
panic("core: nil context provided to ProvideScoped")
|
||||
}
|
||||
if isNil(service) {
|
||||
panic("core: cannot provide nil service")
|
||||
}
|
||||
|
||||
targetType := reflect.TypeFor[T]()
|
||||
targetContainer := ctx.Container()
|
||||
targetContainer.provide(targetType, service)
|
||||
|
||||
+26
-4
@@ -134,6 +134,15 @@ func (c *Context) Parent() *Context {
|
||||
return c.parent
|
||||
}
|
||||
|
||||
// Root returns the root Context in the hierarchy.
|
||||
func (c *Context) Root() *Context {
|
||||
curr := c
|
||||
for curr.parent != nil {
|
||||
curr = curr.parent
|
||||
}
|
||||
return curr
|
||||
}
|
||||
|
||||
// Fork creates a child Context with its own scoped IoC container and values,
|
||||
// linked to this Context for hierarchical fallback resolution and cascading teardown.
|
||||
func (c *Context) Fork() *Context {
|
||||
@@ -337,15 +346,28 @@ func (c *Context) IsDisposed() bool {
|
||||
return c.disposed
|
||||
}
|
||||
|
||||
// RegisterDriver registers a runtime driver engine on this Context.
|
||||
// RegisterDriver registers a runtime driver engine on this Context hierarchy.
|
||||
func (c *Context) RegisterDriver(d Driver) error {
|
||||
if d == nil {
|
||||
return ErrNilService
|
||||
}
|
||||
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
c.drivers = append(c.drivers, d)
|
||||
root := c.Root()
|
||||
root.mu.Lock()
|
||||
root.drivers = append(root.drivers, d)
|
||||
root.mu.Unlock()
|
||||
|
||||
c.OnDispose(func() error {
|
||||
root.mu.Lock()
|
||||
defer root.mu.Unlock()
|
||||
for i, drv := range root.drivers {
|
||||
if drv == d {
|
||||
root.drivers = append(root.drivers[:i], root.drivers[i+1:]...)
|
||||
break
|
||||
}
|
||||
}
|
||||
return nil
|
||||
})
|
||||
return nil
|
||||
}
|
||||
|
||||
|
||||
@@ -238,14 +238,14 @@ func TestContextHierarchyAndFork(t *testing.T) {
|
||||
|
||||
// Child provides LogService
|
||||
childLog := &logServiceImpl{}
|
||||
core.Provide[LogService](child, childLog)
|
||||
core.ProvideScoped[LogService](child, childLog)
|
||||
|
||||
// Child has LogService, parent does not
|
||||
assert.True(t, core.Has[LogService](child))
|
||||
assert.False(t, core.Has[LogService](parent))
|
||||
|
||||
// Child overrides SampleService
|
||||
core.Provide[SampleService](child, &sampleServiceImpl{prefix: "Child:"})
|
||||
// Child overrides SampleService locally
|
||||
core.ProvideScoped[SampleService](child, &sampleServiceImpl{prefix: "Child:"})
|
||||
childSvc, err := core.Inject[SampleService](child)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, "Child: Ryan", childSvc.Greet("Ryan"))
|
||||
|
||||
@@ -0,0 +1,157 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package core
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"reflect"
|
||||
"sync"
|
||||
)
|
||||
|
||||
// FiberState represents the lifecycle status of a plugin instance in the Cordis micro-kernel.
|
||||
type FiberState string
|
||||
|
||||
const (
|
||||
// FiberPending indicates the plugin is waiting for its required dependencies to be provided.
|
||||
FiberPending FiberState = "PENDING"
|
||||
|
||||
// FiberLoading indicates the plugin is currently running its Apply mounting phase.
|
||||
FiberLoading FiberState = "LOADING"
|
||||
|
||||
// FiberActive indicates the plugin is fully mounted, active, and operating without error.
|
||||
FiberActive FiberState = "ACTIVE"
|
||||
|
||||
// FiberUnloading indicates the plugin is tearing down its scoped effects in LIFO order.
|
||||
FiberUnloading FiberState = "UNLOADING"
|
||||
|
||||
// FiberDisposed indicates the plugin has been completely unmounted and its context disposed.
|
||||
FiberDisposed FiberState = "DISPOSED"
|
||||
)
|
||||
|
||||
// Fiber wraps a Plugin instance with a dedicated scoped Context and manages its
|
||||
// reactive lifecycle state machine according to Cordis spatiotemporal composability principles.
|
||||
type Fiber struct {
|
||||
mu sync.RWMutex
|
||||
plugin Plugin
|
||||
state FiberState
|
||||
ctx *Context
|
||||
deps []reflect.Type
|
||||
err error
|
||||
}
|
||||
|
||||
// NewFiber creates a new Fiber for the specified plugin with a child scoped context.
|
||||
func NewFiber(rootCtx *Context, plugin Plugin) *Fiber {
|
||||
var deps []reflect.Type
|
||||
if depPlugin, ok := plugin.(DependentPlugin); ok {
|
||||
deps = depPlugin.Inject()
|
||||
}
|
||||
|
||||
return &Fiber{
|
||||
plugin: plugin,
|
||||
state: FiberPending,
|
||||
ctx: rootCtx.Fork(),
|
||||
deps: deps,
|
||||
}
|
||||
}
|
||||
|
||||
// Plugin returns the underlying Plugin instance.
|
||||
func (f *Fiber) Plugin() Plugin {
|
||||
return f.plugin
|
||||
}
|
||||
|
||||
// Name returns the unique identifier of the plugin.
|
||||
func (f *Fiber) Name() string {
|
||||
if f.plugin == nil {
|
||||
return ""
|
||||
}
|
||||
return f.plugin.Name()
|
||||
}
|
||||
|
||||
// State returns the current lifecycle state of this Fiber.
|
||||
func (f *Fiber) State() FiberState {
|
||||
f.mu.RLock()
|
||||
defer f.mu.RUnlock()
|
||||
return f.state
|
||||
}
|
||||
|
||||
// Context returns the dedicated scoped Context for this Fiber.
|
||||
func (f *Fiber) Context() *Context {
|
||||
return f.ctx
|
||||
}
|
||||
|
||||
// Dependencies returns the list of required service reflect.Types.
|
||||
func (f *Fiber) Dependencies() []reflect.Type {
|
||||
res := make([]reflect.Type, len(f.deps))
|
||||
copy(res, f.deps)
|
||||
return res
|
||||
}
|
||||
|
||||
// Error returns the latest mounting or unmounting error, if any.
|
||||
func (f *Fiber) Error() error {
|
||||
f.mu.RLock()
|
||||
defer f.mu.RUnlock()
|
||||
return f.err
|
||||
}
|
||||
|
||||
// DependenciesSatisfied checks if all declared dependencies are present in the target Context container.
|
||||
func (f *Fiber) DependenciesSatisfied(ctx *Context) bool {
|
||||
if len(f.deps) == 0 {
|
||||
return true
|
||||
}
|
||||
container := ctx.Container()
|
||||
for _, dep := range f.deps {
|
||||
if _, err := container.resolve(dep); err != nil {
|
||||
return false
|
||||
}
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
// Load executes the plugin mounting lifecycle: PENDING -> LOADING -> ACTIVE.
|
||||
func (f *Fiber) Load() error {
|
||||
f.mu.Lock()
|
||||
if f.state != FiberPending {
|
||||
f.mu.Unlock()
|
||||
return nil
|
||||
}
|
||||
f.state = FiberLoading
|
||||
f.err = nil
|
||||
f.mu.Unlock()
|
||||
|
||||
if err := f.plugin.Apply(f.ctx); err != nil {
|
||||
f.mu.Lock()
|
||||
f.err = fmt.Errorf("fiber %q: apply failed: %w", f.Name(), err)
|
||||
f.state = FiberPending
|
||||
_ = f.ctx.Dispose()
|
||||
f.mu.Unlock()
|
||||
return err
|
||||
}
|
||||
|
||||
f.mu.Lock()
|
||||
f.state = FiberActive
|
||||
f.mu.Unlock()
|
||||
return nil
|
||||
}
|
||||
|
||||
// Unload tears down the plugin: ACTIVE -> UNLOADING -> DISPOSED.
|
||||
func (f *Fiber) Unload() error {
|
||||
f.mu.Lock()
|
||||
if f.state != FiberActive && f.state != FiberLoading {
|
||||
f.mu.Unlock()
|
||||
return nil
|
||||
}
|
||||
f.state = FiberUnloading
|
||||
f.mu.Unlock()
|
||||
|
||||
err := f.ctx.Dispose()
|
||||
|
||||
f.mu.Lock()
|
||||
f.state = FiberDisposed
|
||||
if err != nil {
|
||||
f.err = err
|
||||
}
|
||||
f.mu.Unlock()
|
||||
|
||||
return err
|
||||
}
|
||||
@@ -0,0 +1,111 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package core_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"reflect"
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
|
||||
"Wavelet/core"
|
||||
)
|
||||
|
||||
type MockServiceA interface {
|
||||
DoA() string
|
||||
}
|
||||
|
||||
type mockServiceAImpl struct{}
|
||||
|
||||
func (m *mockServiceAImpl) DoA() string { return "doneA" }
|
||||
|
||||
type MockServiceB interface {
|
||||
DoB() string
|
||||
}
|
||||
|
||||
type mockServiceBImpl struct{}
|
||||
|
||||
func (m *mockServiceBImpl) DoB() string { return "doneB" }
|
||||
|
||||
type mockProviderPlugin struct {
|
||||
name string
|
||||
applied bool
|
||||
}
|
||||
|
||||
func (p *mockProviderPlugin) Name() string {
|
||||
return p.name
|
||||
}
|
||||
|
||||
func (p *mockProviderPlugin) Apply(ctx *core.Context) error {
|
||||
p.applied = true
|
||||
core.Provide[MockServiceA](ctx, &mockServiceAImpl{})
|
||||
return nil
|
||||
}
|
||||
|
||||
type mockConsumerPlugin struct {
|
||||
name string
|
||||
applied bool
|
||||
gotSvc MockServiceA
|
||||
}
|
||||
|
||||
func (p *mockConsumerPlugin) Name() string {
|
||||
return p.name
|
||||
}
|
||||
|
||||
func (p *mockConsumerPlugin) Inject() []reflect.Type {
|
||||
return []reflect.Type{
|
||||
reflect.TypeFor[MockServiceA](),
|
||||
}
|
||||
}
|
||||
|
||||
func (p *mockConsumerPlugin) Apply(ctx *core.Context) error {
|
||||
p.applied = true
|
||||
svc, err := core.Inject[MockServiceA](ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
p.gotSvc = svc
|
||||
return nil
|
||||
}
|
||||
|
||||
func TestFiber_ConfluenceAndReactiveActivation(t *testing.T) {
|
||||
app := core.NewApp()
|
||||
|
||||
// Register Consumer BEFORE Provider to test confluence & out-of-order dependency resolution
|
||||
consumer := &mockConsumerPlugin{name: "consumer-plugin"}
|
||||
provider := &mockProviderPlugin{name: "provider-plugin"}
|
||||
|
||||
app.Use(consumer, provider)
|
||||
|
||||
err := app.Start(context.Background())
|
||||
require.NoError(t, err)
|
||||
|
||||
assert.True(t, provider.applied, "provider should be applied")
|
||||
assert.True(t, consumer.applied, "consumer should be reactively applied once dependency was provided")
|
||||
assert.NotNil(t, consumer.gotSvc)
|
||||
assert.Equal(t, "doneA", consumer.gotSvc.DoA())
|
||||
|
||||
fibers := app.Fibers()
|
||||
require.Equal(t, 2, len(fibers))
|
||||
for _, f := range fibers {
|
||||
assert.Equal(t, core.FiberActive, f.State())
|
||||
}
|
||||
|
||||
err = app.Stop()
|
||||
assert.NoError(t, err)
|
||||
}
|
||||
|
||||
func TestFiber_UnsatisfiedDependencyReturnsError(t *testing.T) {
|
||||
app := core.NewApp()
|
||||
|
||||
// Register Consumer whose dependency is never provided
|
||||
consumer := &mockConsumerPlugin{name: "consumer-plugin"}
|
||||
app.Use(consumer)
|
||||
|
||||
err := app.Start(context.Background())
|
||||
assert.Error(t, err)
|
||||
assert.Contains(t, err.Error(), "unsatisfied dependencies")
|
||||
}
|
||||
@@ -6,6 +6,7 @@ package core
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"reflect"
|
||||
|
||||
"Wavelet/core/extpoints"
|
||||
)
|
||||
@@ -48,6 +49,12 @@ type Plugin interface {
|
||||
Apply(ctx *Context) error
|
||||
}
|
||||
|
||||
// DependentPlugin is an optional extension interface for plugins that declare required service dependencies.
|
||||
type DependentPlugin interface {
|
||||
Plugin
|
||||
Inject() []reflect.Type
|
||||
}
|
||||
|
||||
// PluginWithManifest is an optional extension interface for plugins that declare metadata.
|
||||
type PluginWithManifest interface {
|
||||
Plugin
|
||||
|
||||
Reference in New Issue
Block a user