fix: persist browser captures atomically
This commit is contained in:
@@ -57,6 +57,8 @@ type WorkspaceProvider func() string
|
||||
type Options struct {
|
||||
RequireToken bool
|
||||
ReceiverToken string
|
||||
Available func() bool
|
||||
Persist func(events.Event) error
|
||||
}
|
||||
|
||||
type Server struct {
|
||||
@@ -128,6 +130,17 @@ func (r *Receiver) SetReceiverToken(token string) {
|
||||
r.options.ReceiverToken = strings.TrimSpace(token)
|
||||
}
|
||||
|
||||
// SetPersistence configures durable capture storage and its availability gate.
|
||||
func (r *Receiver) SetPersistence(available func() bool, persist func(events.Event) error) {
|
||||
if r == nil {
|
||||
return
|
||||
}
|
||||
r.optionsMu.Lock()
|
||||
defer r.optionsMu.Unlock()
|
||||
r.options.Available = available
|
||||
r.options.Persist = persist
|
||||
}
|
||||
|
||||
func Start(addr string, receiver *Receiver) (*Server, error) {
|
||||
if receiver == nil {
|
||||
return nil, fmt.Errorf("receiver is required")
|
||||
@@ -181,6 +194,11 @@ func (r *Receiver) ServeHTTP(w http.ResponseWriter, req *http.Request) {
|
||||
writeError(w, http.StatusUnauthorized, err.Error())
|
||||
return
|
||||
}
|
||||
options := r.currentOptions()
|
||||
if options.Available != nil && !options.Available() {
|
||||
writeError(w, http.StatusServiceUnavailable, "browser inbox unavailable")
|
||||
return
|
||||
}
|
||||
|
||||
defer req.Body.Close()
|
||||
decoder := json.NewDecoder(http.MaxBytesReader(w, req.Body, maxCaptureBodyBytes))
|
||||
@@ -209,17 +227,27 @@ func (r *Receiver) ServeHTTP(w http.ResponseWriter, req *http.Request) {
|
||||
}
|
||||
|
||||
eventName := "browser.capture." + payload.Kind
|
||||
if r.bus == nil || !r.bus.HasSubscribers(eventName) {
|
||||
if options.Persist == nil && (r.bus == nil || !r.bus.HasSubscribers(eventName)) {
|
||||
writeError(w, http.StatusServiceUnavailable, "browser inbox unavailable")
|
||||
return
|
||||
}
|
||||
eventPayload := payload.EventPayload()
|
||||
r.annotateWorkspace(eventPayload)
|
||||
r.bus.Publish(events.Event{
|
||||
event := events.Event{
|
||||
Name: eventName,
|
||||
Timestamp: time.Now().UTC().Format(time.RFC3339Nano),
|
||||
Payload: eventPayload,
|
||||
})
|
||||
}
|
||||
if options.Persist != nil {
|
||||
if err := options.Persist(event); err != nil {
|
||||
log.Printf("[browserreceiver] persist %s: %v", payload.CaptureID, err)
|
||||
writeError(w, http.StatusServiceUnavailable, "browser inbox unavailable")
|
||||
return
|
||||
}
|
||||
}
|
||||
if r.bus != nil {
|
||||
r.bus.Publish(event)
|
||||
}
|
||||
|
||||
w.WriteHeader(http.StatusAccepted)
|
||||
_ = json.NewEncoder(w).Encode(map[string]string{
|
||||
@@ -228,6 +256,15 @@ func (r *Receiver) ServeHTTP(w http.ResponseWriter, req *http.Request) {
|
||||
})
|
||||
}
|
||||
|
||||
func (r *Receiver) currentOptions() Options {
|
||||
if r == nil {
|
||||
return Options{}
|
||||
}
|
||||
r.optionsMu.RLock()
|
||||
defer r.optionsMu.RUnlock()
|
||||
return r.options
|
||||
}
|
||||
|
||||
func (r *Receiver) validateReceiverToken(req *http.Request) error {
|
||||
if r == nil {
|
||||
return nil
|
||||
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"bytes"
|
||||
"encoding/base64"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
@@ -321,6 +322,61 @@ func TestReceiverAcceptsPairedToken(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestReceiverRejectsCaptureWhenInboxIsUnavailable(t *testing.T) {
|
||||
bus := events.NewBus()
|
||||
bus.Subscribe("browser.capture.page", func(event events.Event) {
|
||||
t.Fatalf("unexpected event while inbox is unavailable: %#v", event)
|
||||
})
|
||||
receiver := NewWithOptions(bus, Options{
|
||||
Available: func() bool { return false },
|
||||
})
|
||||
req := httptest.NewRequest(http.MethodPost, capturePath, strings.NewReader(`{
|
||||
"schemaVersion": 1,
|
||||
"captureId": "capture-no-vault",
|
||||
"capturedAt": "2026-07-11T12:00:00Z",
|
||||
"kind": "page",
|
||||
"page": {"url": "https://example.com"}
|
||||
}`))
|
||||
res := httptest.NewRecorder()
|
||||
|
||||
receiver.ServeHTTP(res, req)
|
||||
|
||||
if res.Code != http.StatusServiceUnavailable {
|
||||
t.Fatalf("status = %d, want %d; body=%s", res.Code, http.StatusServiceUnavailable, res.Body.String())
|
||||
}
|
||||
}
|
||||
|
||||
func TestReceiverDoesNotAcknowledgeFailedPersistence(t *testing.T) {
|
||||
bus := events.NewBus()
|
||||
published := false
|
||||
bus.Subscribe("browser.capture.page", func(event events.Event) {
|
||||
published = true
|
||||
})
|
||||
receiver := NewWithOptions(bus, Options{
|
||||
Available: func() bool { return true },
|
||||
Persist: func(event events.Event) error {
|
||||
return errors.New("disk full")
|
||||
},
|
||||
})
|
||||
req := httptest.NewRequest(http.MethodPost, capturePath, strings.NewReader(`{
|
||||
"schemaVersion": 1,
|
||||
"captureId": "capture-disk-full",
|
||||
"capturedAt": "2026-07-11T12:00:00Z",
|
||||
"kind": "page",
|
||||
"page": {"url": "https://example.com"}
|
||||
}`))
|
||||
res := httptest.NewRecorder()
|
||||
|
||||
receiver.ServeHTTP(res, req)
|
||||
|
||||
if res.Code != http.StatusServiceUnavailable {
|
||||
t.Fatalf("status = %d, want %d; body=%s", res.Code, http.StatusServiceUnavailable, res.Body.String())
|
||||
}
|
||||
if published {
|
||||
t.Fatal("capture event was published after persistence failed")
|
||||
}
|
||||
}
|
||||
|
||||
func TestReceiverRotatesPairedToken(t *testing.T) {
|
||||
bus := events.NewBus()
|
||||
bus.Subscribe("browser.capture.page", func(event events.Event) {})
|
||||
|
||||
@@ -84,13 +84,22 @@ func atomicWrite(path string, data []byte) error {
|
||||
// ReadPluginSettings reads all settings for a plugin.
|
||||
// Returns empty map if settings.json does not exist.
|
||||
func (s *Storage) ReadPluginSettings(pluginID string) (map[string]interface{}, error) {
|
||||
s.mu.RLock()
|
||||
defer s.mu.RUnlock()
|
||||
if err := validatePluginID(pluginID); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
var result map[string]interface{}
|
||||
err := s.withOpenVault(func(vaultPath string) error {
|
||||
var err error
|
||||
result, err = readPluginSettingsAt(vaultPath, pluginID)
|
||||
return err
|
||||
})
|
||||
return result, err
|
||||
}
|
||||
|
||||
dir := s.vault.GetPluginSettingsPath(pluginID)
|
||||
path := filepath.Join(dir, "settings.json")
|
||||
|
||||
func readPluginSettingsAt(vaultPath, pluginID string) (map[string]interface{}, error) {
|
||||
path := filepath.Join(vaultPath, ".verstak", "plugin-settings", pluginID, "settings.json")
|
||||
data, err := os.ReadFile(path)
|
||||
if err != nil {
|
||||
if os.IsNotExist(err) {
|
||||
@@ -108,13 +117,18 @@ func (s *Storage) ReadPluginSettings(pluginID string) (map[string]interface{}, e
|
||||
|
||||
// WritePluginSettings writes all settings for a plugin atomically.
|
||||
func (s *Storage) WritePluginSettings(pluginID string, data map[string]interface{}) error {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
if err := validatePluginID(pluginID); err != nil {
|
||||
return err
|
||||
}
|
||||
return s.withOpenVault(func(vaultPath string) error {
|
||||
return writePluginSettingsAt(vaultPath, pluginID, data)
|
||||
})
|
||||
}
|
||||
|
||||
dir := s.vault.GetPluginSettingsPath(pluginID)
|
||||
path := filepath.Join(dir, "settings.json")
|
||||
|
||||
func writePluginSettingsAt(vaultPath, pluginID string, data map[string]interface{}) error {
|
||||
path := filepath.Join(vaultPath, ".verstak", "plugin-settings", pluginID, "settings.json")
|
||||
encoded, err := json.MarshalIndent(data, "", " ")
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to marshal settings for plugin %s: %w", pluginID, err)
|
||||
@@ -137,12 +151,39 @@ func (s *Storage) ReadPluginSetting(pluginID, key string) (interface{}, error) {
|
||||
|
||||
// WritePluginSetting writes a single setting key.
|
||||
func (s *Storage) WritePluginSetting(pluginID, key string, value interface{}) error {
|
||||
settings, err := s.ReadPluginSettings(pluginID)
|
||||
if err != nil {
|
||||
return s.UpdatePluginSettings(pluginID, func(settings map[string]interface{}) error {
|
||||
settings[key] = value
|
||||
return nil
|
||||
})
|
||||
}
|
||||
|
||||
// UpdatePluginSettings atomically reads, mutates, and writes a plugin's settings.
|
||||
func (s *Storage) UpdatePluginSettings(pluginID string, update func(map[string]interface{}) error) error {
|
||||
if update == nil {
|
||||
return fmt.Errorf("settings update is nil")
|
||||
}
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
if err := validatePluginID(pluginID); err != nil {
|
||||
return err
|
||||
}
|
||||
settings[key] = value
|
||||
return s.WritePluginSettings(pluginID, settings)
|
||||
return s.withOpenVault(func(vaultPath string) error {
|
||||
settings, err := readPluginSettingsAt(vaultPath, pluginID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := update(settings); err != nil {
|
||||
return err
|
||||
}
|
||||
return writePluginSettingsAt(vaultPath, pluginID, settings)
|
||||
})
|
||||
}
|
||||
|
||||
func (s *Storage) withOpenVault(operation func(string) error) error {
|
||||
if s == nil || s.vault == nil {
|
||||
return fmt.Errorf("vault is not initialized")
|
||||
}
|
||||
return s.vault.WithOpenPath(operation)
|
||||
}
|
||||
|
||||
// ─── Data JSON API ────────────────────────────────────────
|
||||
|
||||
@@ -70,6 +70,23 @@ func (v *Vault) GetVaultPath() string {
|
||||
return v.path
|
||||
}
|
||||
|
||||
// WithOpenPath runs fn while holding a read lease on the currently open vault.
|
||||
// CloseVault and OpenVault wait until the operation completes.
|
||||
func (v *Vault) WithOpenPath(fn func(string) error) error {
|
||||
if v == nil {
|
||||
return fmt.Errorf("vault is not initialized")
|
||||
}
|
||||
if fn == nil {
|
||||
return fmt.Errorf("vault path operation is nil")
|
||||
}
|
||||
v.mu.RLock()
|
||||
defer v.mu.RUnlock()
|
||||
if v.status != StatusOpen || strings.TrimSpace(v.path) == "" {
|
||||
return fmt.Errorf("vault is not open")
|
||||
}
|
||||
return fn(v.path)
|
||||
}
|
||||
|
||||
// GetVaultMeta returns the current vault metadata.
|
||||
func (v *Vault) GetVaultMeta() *VaultMeta {
|
||||
v.mu.RLock()
|
||||
|
||||
Reference in New Issue
Block a user