Skip to content

Commit c8d314e

Browse files
committed
fix: PluginData/PluginStore return errors, idempotent Close, schemaCache cleanup, Values hydration
1 parent 626c3c0 commit c8d314e

7 files changed

Lines changed: 93 additions & 64 deletions

File tree

backend/pkg/plugin/data/controller.go

Lines changed: 18 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -24,12 +24,12 @@ var _ Controller = (*controller)(nil)
2424

2525
type controller struct {
2626
logger logging.Logger
27-
pluginDataFn func(pluginID string) *appstate.ScopedRoot
27+
pluginDataFn func(pluginID string) (*appstate.ScopedRoot, error)
2828
}
2929

3030
// NewController creates a new data store controller.
3131
// pluginDataFn returns a ScopedRoot for the given plugin's data directory.
32-
func NewController(logger logging.Logger, pluginDataFn func(string) *appstate.ScopedRoot) Controller {
32+
func NewController(logger logging.Logger, pluginDataFn func(string) (*appstate.ScopedRoot, error)) Controller {
3333
return &controller{
3434
logger: logger.Named("DataController"),
3535
pluginDataFn: pluginDataFn,
@@ -39,7 +39,10 @@ func NewController(logger logging.Logger, pluginDataFn func(string) *appstate.Sc
3939
func (c *controller) Get(pluginID, key string) (any, error) {
4040
logger := c.logger.With(logging.Any("pluginID", pluginID), logging.Any("key", key))
4141

42-
root := c.pluginDataFn(pluginID)
42+
root, err := c.pluginDataFn(pluginID)
43+
if err != nil {
44+
return nil, err
45+
}
4346
data, err := root.ReadFile(key + ".json")
4447
if err != nil {
4548
if errors.Is(err, os.ErrNotExist) {
@@ -61,7 +64,10 @@ func (c *controller) Get(pluginID, key string) (any, error) {
6164
func (c *controller) Set(pluginID, key string, value any) error {
6265
logger := c.logger.With(logging.Any("pluginID", pluginID), logging.Any("key", key))
6366

64-
root := c.pluginDataFn(pluginID)
67+
root, err := c.pluginDataFn(pluginID)
68+
if err != nil {
69+
return err
70+
}
6571
data, err := json.MarshalIndent(value, "", " ")
6672
if err != nil {
6773
logger.Errorw(context.Background(), "failed to marshal data", "error", err)
@@ -79,7 +85,10 @@ func (c *controller) Set(pluginID, key string, value any) error {
7985
func (c *controller) Delete(pluginID, key string) error {
8086
logger := c.logger.With(logging.Any("pluginID", pluginID), logging.Any("key", key))
8187

82-
root := c.pluginDataFn(pluginID)
88+
root, err := c.pluginDataFn(pluginID)
89+
if err != nil {
90+
return err
91+
}
8392
if err := root.Remove(key + ".json"); err != nil {
8493
if errors.Is(err, os.ErrNotExist) {
8594
return nil
@@ -94,7 +103,10 @@ func (c *controller) Delete(pluginID, key string) error {
94103
func (c *controller) Keys(pluginID string) ([]string, error) {
95104
logger := c.logger.With(logging.Any("pluginID", pluginID))
96105

97-
root := c.pluginDataFn(pluginID)
106+
root, err := c.pluginDataFn(pluginID)
107+
if err != nil {
108+
return nil, err
109+
}
98110
entries, err := root.ReadDir(".")
99111
if err != nil {
100112
if errors.Is(err, os.ErrNotExist) {

backend/pkg/plugin/resource/controller.go

Lines changed: 18 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -78,7 +78,7 @@ type controller struct {
7878
graph *graph.RelationshipGraph
7979

8080
onCrashCallback func(pluginID string)
81-
pluginStoreFn func(pluginID string) *appstate.ScopedRoot
81+
pluginStoreFn func(pluginID string) (*appstate.ScopedRoot, error)
8282
}
8383

8484
// compile-time assertions
@@ -89,7 +89,7 @@ var (
8989

9090
// NewController creates a new resource Controller.
9191
// pluginStoreFn returns a ScopedRoot for the given plugin's store directory.
92-
func NewController(logger logging.Logger, sp pkgsettings.Provider, pluginStoreFn func(string) *appstate.ScopedRoot) Controller {
92+
func NewController(logger logging.Logger, sp pkgsettings.Provider, pluginStoreFn func(string) (*appstate.ScopedRoot, error)) Controller {
9393
store := registry.NewMemoryStore()
9494
g := graph.NewRelationshipGraph()
9595
graphIndexer := graph.NewGraphIndexer(g, store)
@@ -163,8 +163,11 @@ func (c *controller) OnPluginInit(pluginID string, meta config.PluginMeta) {
163163
logger.Debugw(context.Background(), "OnPluginInit")
164164

165165
// Load persisted connections from disk.
166-
state, err := loadFromLocalStore(c.pluginStoreFn(pluginID))
167-
if err != nil {
166+
var state map[string][]types.Connection
167+
if storeRoot, err := c.pluginStoreFn(pluginID); err != nil {
168+
logger.Errorw(context.Background(), "failed to resolve plugin store root", "error", err)
169+
state = make(map[string][]types.Connection)
170+
} else if state, err = loadFromLocalStore(storeRoot); err != nil {
168171
logger.Errorw(context.Background(), "failed to load connections from local store", "error", err)
169172
state = make(map[string][]types.Connection)
170173
}
@@ -268,7 +271,9 @@ func (c *controller) OnPluginStop(pluginID string, meta config.PluginMeta) error
268271
c.connsMu.RLock()
269272
conns := c.connections
270273
c.connsMu.RUnlock()
271-
if err := saveToLocalStore(c.pluginStoreFn(pluginID), conns); err != nil {
274+
if storeRoot, err := c.pluginStoreFn(pluginID); err != nil {
275+
logger.Errorw(context.Background(), "failed to resolve plugin store root", "error", err)
276+
} else if err := saveToLocalStore(storeRoot, conns); err != nil {
272277
logger.Errorw(context.Background(), "failed to save connections to local store", "error", err)
273278
}
274279

@@ -322,7 +327,9 @@ func (c *controller) OnPluginShutdown(pluginID string, meta config.PluginMeta) e
322327
func (c *controller) OnPluginDestroy(pluginID string, meta config.PluginMeta) error {
323328
logger := c.logger.With(logging.Any("pluginID", pluginID))
324329
logger.Debugw(context.Background(), "OnPluginDestroy")
325-
if err := removeLocalStore(c.pluginStoreFn(pluginID)); err != nil {
330+
if storeRoot, err := c.pluginStoreFn(pluginID); err != nil {
331+
logger.Errorw(context.Background(), "failed to resolve plugin store root", "error", err)
332+
} else if err := removeLocalStore(storeRoot); err != nil {
326333
logger.Errorw(context.Background(), "failed to remove local store", "error", err)
327334
}
328335
return nil
@@ -638,9 +645,11 @@ func (c *controller) LoadConnections(pluginID string) ([]types.Connection, error
638645
c.connsMu.Unlock()
639646

640647
// Best-effort persist.
641-
c.connsMu.RLock()
642-
_ = saveToLocalStore(c.pluginStoreFn(pluginID), c.connections)
643-
c.connsMu.RUnlock()
648+
if storeRoot, err := c.pluginStoreFn(pluginID); err == nil {
649+
c.connsMu.RLock()
650+
_ = saveToLocalStore(storeRoot, c.connections)
651+
c.connsMu.RUnlock()
652+
}
644653

645654
return conns, nil
646655
}

backend/pkg/plugin/resource/controller_apperror_test.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,7 @@ import (
2020
func newTestController() *controller {
2121
// Provide a pluginStoreFn that panics if called — these tests never exercise
2222
// the local store path, so the function should never be invoked.
23-
storeFn := func(pluginID string) *appstate.ScopedRoot {
23+
storeFn := func(pluginID string) (*appstate.ScopedRoot, error) {
2424
panic("unexpected call to pluginStoreFn in test for plugin " + pluginID)
2525
}
2626
ctrl := NewController(logging.NewNop(), nil, storeFn).(*controller)

backend/pkg/plugin/resource/testutil_test.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -600,7 +600,7 @@ func newTestControllerWithEmitter(t *testing.T) (*controller, *recordingEmitter)
600600
registryStore: store,
601601
dispatcher: disp,
602602
graph: g,
603-
pluginStoreFn: svc.PluginStore,
603+
pluginStoreFn: svc.PluginStore,
604604
}
605605
return ctrl, emitter
606606
}

backend/pkg/plugin/settings/controller.go

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -345,6 +345,8 @@ func (c *controller) OnPluginDestroy(pluginID string, meta config.PluginMeta) er
345345

346346
c.mu.Lock()
347347
delete(c.hydrated, pluginID)
348+
delete(c.clients, pluginID)
349+
delete(c.schemaCache, pluginID)
348350
c.mu.Unlock()
349351

350352
return nil
@@ -423,6 +425,9 @@ func (c *controller) Values() map[string]any {
423425

424426
values := make(map[string]any)
425427
for pluginID, client := range snapshot {
428+
// Attempt hydration if not yet done
429+
c.tryHydrate(ctx, pluginID)
430+
426431
clientValues := client.ListSettings()
427432
if clientValues == nil {
428433
// gRPC failed, fall back to bbolt

internal/appstate/appstate.go

Lines changed: 46 additions & 45 deletions
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,9 @@ type Service struct {
4141
pluginRootsMu sync.Mutex
4242
pluginDataRoots map[string]*ScopedRoot
4343
pluginStoreRoots map[string]*ScopedRoot
44+
45+
closeOnce sync.Once
46+
closeErr error
4447
}
4548

4649
type config struct {
@@ -140,41 +143,45 @@ func (s *Service) closePartial() {
140143

141144
// Close releases all resources held by the Service.
142145
// Errors from individual close operations are aggregated and returned together.
146+
// Close is idempotent; subsequent calls return the result of the first call.
143147
func (s *Service) Close() error {
144-
var allErrors []error
145-
146-
s.pluginRootsMu.Lock()
147-
for id, r := range s.pluginDataRoots {
148-
if err := r.Close(); err != nil {
149-
allErrors = append(allErrors, fmt.Errorf("close plugin data root %q: %w", id, err))
148+
s.closeOnce.Do(func() {
149+
var allErrors []error
150+
151+
s.pluginRootsMu.Lock()
152+
for id, r := range s.pluginDataRoots {
153+
if err := r.Close(); err != nil {
154+
allErrors = append(allErrors, fmt.Errorf("close plugin data root %q: %w", id, err))
155+
}
150156
}
151-
}
152-
for id, r := range s.pluginStoreRoots {
153-
if err := r.Close(); err != nil {
154-
allErrors = append(allErrors, fmt.Errorf("close plugin store root %q: %w", id, err))
157+
for id, r := range s.pluginStoreRoots {
158+
if err := r.Close(); err != nil {
159+
allErrors = append(allErrors, fmt.Errorf("close plugin store root %q: %w", id, err))
160+
}
155161
}
156-
}
157-
s.pluginDataRoots = nil
158-
s.pluginStoreRoots = nil
159-
s.pluginRootsMu.Unlock()
162+
s.pluginDataRoots = nil
163+
s.pluginStoreRoots = nil
164+
s.pluginRootsMu.Unlock()
160165

161-
if err := s.plugins.Close(); err != nil {
162-
allErrors = append(allErrors, fmt.Errorf("close plugins root: %w", err))
163-
}
164-
if err := s.logs.Close(); err != nil {
165-
allErrors = append(allErrors, fmt.Errorf("close logs root: %w", err))
166-
}
167-
if err := s.rootDir.Close(); err != nil {
168-
allErrors = append(allErrors, fmt.Errorf("close root dir: %w", err))
169-
}
170-
if s.osRoot != nil {
171-
if err := s.osRoot.Close(); err != nil {
172-
allErrors = append(allErrors, fmt.Errorf("close os root: %w", err))
166+
if err := s.plugins.Close(); err != nil {
167+
allErrors = append(allErrors, fmt.Errorf("close plugins root: %w", err))
173168
}
174-
}
175-
releaseLock(s.lock)
169+
if err := s.logs.Close(); err != nil {
170+
allErrors = append(allErrors, fmt.Errorf("close logs root: %w", err))
171+
}
172+
if err := s.rootDir.Close(); err != nil {
173+
allErrors = append(allErrors, fmt.Errorf("close root dir: %w", err))
174+
}
175+
if s.osRoot != nil {
176+
if err := s.osRoot.Close(); err != nil {
177+
allErrors = append(allErrors, fmt.Errorf("close os root: %w", err))
178+
}
179+
}
180+
releaseLock(s.lock)
176181

177-
return errors.Join(allErrors...)
182+
s.closeErr = errors.Join(allErrors...)
183+
})
184+
return s.closeErr
178185
}
179186

180187
// Root returns the absolute path to the state directory.
@@ -191,52 +198,46 @@ func (s *Service) RootDir() *ScopedRoot { return s.rootDir }
191198

192199
// PluginData returns a scoped root for a specific plugin's data directory.
193200
// The result is cached; subsequent calls for the same id return the cached instance.
194-
// Returns nil if the scoped root cannot be created.
195-
func (s *Service) PluginData(id string) *ScopedRoot {
201+
func (s *Service) PluginData(id string) (*ScopedRoot, error) {
196202
if err := validatePluginID(id); err != nil {
197-
fmt.Fprintf(os.Stderr, "appstate: PluginData: %v\n", err)
198-
return nil
203+
return nil, fmt.Errorf("appstate: PluginData: %w", err)
199204
}
200205

201206
s.pluginRootsMu.Lock()
202207
defer s.pluginRootsMu.Unlock()
203208

204209
if r, ok := s.pluginDataRoots[id]; ok {
205-
return r
210+
return r, nil
206211
}
207212

208213
dir := filepath.Join(s.root, "plugins", id, "data")
209214
r, err := newScopedRoot(dir)
210215
if err != nil {
211-
fmt.Fprintf(os.Stderr, "appstate: failed to create plugin data root for %q: %v\n", id, err)
212-
return nil
216+
return nil, fmt.Errorf("appstate: create plugin data root for %q: %w", id, err)
213217
}
214218
s.pluginDataRoots[id] = r
215-
return r
219+
return r, nil
216220
}
217221

218222
// PluginStore returns a scoped root for a specific plugin's store directory.
219223
// The result is cached; subsequent calls for the same id return the cached instance.
220-
// Returns nil if the scoped root cannot be created.
221-
func (s *Service) PluginStore(id string) *ScopedRoot {
224+
func (s *Service) PluginStore(id string) (*ScopedRoot, error) {
222225
if err := validatePluginID(id); err != nil {
223-
fmt.Fprintf(os.Stderr, "appstate: PluginStore: %v\n", err)
224-
return nil
226+
return nil, fmt.Errorf("appstate: PluginStore: %w", err)
225227
}
226228

227229
s.pluginRootsMu.Lock()
228230
defer s.pluginRootsMu.Unlock()
229231

230232
if r, ok := s.pluginStoreRoots[id]; ok {
231-
return r
233+
return r, nil
232234
}
233235

234236
dir := filepath.Join(s.root, "plugins", id, "store")
235237
r, err := newScopedRoot(dir)
236238
if err != nil {
237-
fmt.Fprintf(os.Stderr, "appstate: failed to create plugin store root for %q: %v\n", id, err)
238-
return nil
239+
return nil, fmt.Errorf("appstate: create plugin store root for %q: %w", id, err)
239240
}
240241
s.pluginStoreRoots[id] = r
241-
return r
242+
return r, nil
242243
}

internal/appstate/appstate_test.go

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -33,7 +33,8 @@ func TestNew_PluginDataAccessor(t *testing.T) {
3333
svc, err := New(WithRoot(dir), WithFlock(false))
3434
require.NoError(t, err)
3535
defer svc.Close()
36-
pd := svc.PluginData("my-plugin")
36+
pd, err := svc.PluginData("my-plugin")
37+
require.NoError(t, err)
3738
require.NoError(t, pd.WriteFile("key.json", []byte(`{"v":1}`), 0600))
3839
got, err := pd.ReadFile("key.json")
3940
require.NoError(t, err)
@@ -46,7 +47,8 @@ func TestNew_PluginStoreAccessor(t *testing.T) {
4647
svc, err := New(WithRoot(dir), WithFlock(false))
4748
require.NoError(t, err)
4849
defer svc.Close()
49-
ps := svc.PluginStore("my-plugin")
50+
ps, err := svc.PluginStore("my-plugin")
51+
require.NoError(t, err)
5052
require.NoError(t, ps.WriteFile("resource", []byte("binary-data"), 0600))
5153
}
5254

0 commit comments

Comments
 (0)