From 9eed1e11c5f049a08db58ec01814477c9075c356 Mon Sep 17 00:00:00 2001 From: Tobias Gesellchen Date: Sun, 23 Aug 2026 13:42:09 +0200 Subject: [PATCH] fix(marge): route preset/recent/source read-modify-write through Mutate* Converts the remaining GetX-then-SaveX call sites (UpdatePreset, RemovePreset, AddRecent's recent + learned-source persistence, AddSource) to the new datastore.Mutate{Presets,Recents,ConfiguredSources} helpers, closing the lost-update race for good on the actual write path the speaker hits on every preset/recent store. Adds a regression test that fires 6 concurrent UpdatePreset calls (same shape as #614's rapid-fire repro) and asserts none are lost. Verified it reliably fails against the pre-fix code (consistently drops presets across repeated runs) and passes reliably with the fix, including under -race. Co-Authored-By: Claude Sonnet 5 --- .../marge/concurrent_update_preset_test.go | 88 ++++++++++ pkg/service/marge/marge.go | 151 +++++++++--------- 2 files changed, 165 insertions(+), 74 deletions(-) create mode 100644 pkg/service/marge/concurrent_update_preset_test.go diff --git a/pkg/service/marge/concurrent_update_preset_test.go b/pkg/service/marge/concurrent_update_preset_test.go new file mode 100644 index 0000000..ca55c4c --- /dev/null +++ b/pkg/service/marge/concurrent_update_preset_test.go @@ -0,0 +1,88 @@ +package marge + +import ( + "fmt" + "os" + "sync" + "testing" + + "github.com/gesellix/bose-soundtouch/pkg/service/datastore" +) + +// TestConcurrentUpdatePresetNoLostUpdates is a regression test for #614's +// 2026-08-23 reproduction: a reporter's script stored six presets via rapid, +// overlapping PUT .../preset/N requests (visible in the speaker's own log as +// interleaved connection IDs, never waiting for one PUT to complete before +// firing the next). One preset silently vanished from Presets.xml. +// +// UpdatePreset used to do GetPresets, mutate one slot, SavePresets as three +// separate steps with no lock spanning them — a classic lost-update race: +// two concurrent calls can each read the same starting list, mutate +// different slots, and the second writer's SavePresets clobbers the first +// writer's update. Fixed by routing the write through +// datastore.MutatePresets, which holds a single write lock for the whole +// read-mutate-write cycle. +func TestConcurrentUpdatePresetNoLostUpdates(t *testing.T) { + tempDir, err := os.MkdirTemp("", "marge-concurrent-update-preset-*") + if err != nil { + t.Fatalf("tempdir: %v", err) + } + defer func() { _ = os.RemoveAll(tempDir) }() + + ds := datastore.NewDataStore(tempDir) + + account := "1234567" + device := "B0D5CC25479C" + + const presetCount = 6 + + var wg sync.WaitGroup + + errs := make([]error, presetCount) + + for i := 1; i <= presetCount; i++ { + wg.Add(1) + + go func(presetNumber int) { + defer wg.Done() + + // sourceid 10003 is the canonical LOCAL_INTERNET_RADIO built-in + // (see CanonicalSourceByID) — same shape as Henri's own repro + // script, which stored six LOCAL_INTERNET_RADIO presets. + putXML := []byte(fmt.Sprintf(` + + Station %d + 10003 + /custom/v1/playback/station%d + stationurl +`, presetNumber, presetNumber)) + + _, err := UpdatePreset(ds, account, device, presetNumber, putXML) + errs[presetNumber-1] = err + }(i) + } + + wg.Wait() + + for i, err := range errs { + if err != nil { + t.Fatalf("UpdatePreset(preset=%d) returned error: %v", i+1, err) + } + } + + presets, err := ds.GetPresets(account, device) + if err != nil { + t.Fatalf("GetPresets: %v", err) + } + + if len(presets) != presetCount { + t.Fatalf("expected %d presets after %d concurrent UpdatePreset calls, got %d: %+v", presetCount, presetCount, len(presets), presets) + } + + for i, p := range presets { + want := fmt.Sprintf("Station %d", i+1) + if p.Name != want { + t.Errorf("preset slot %d: expected name %q, got %q — a concurrent update was lost", i+1, want, p.Name) + } + } +} diff --git a/pkg/service/marge/marge.go b/pkg/service/marge/marge.go index 2db9ad7..67ded8b 100644 --- a/pkg/service/marge/marge.go +++ b/pkg/service/marge/marge.go @@ -1508,19 +1508,18 @@ func AccountFullToXML(ds *datastore.DataStore, account string) ([]byte, error) { // RemovePreset clears a preset for the specified account and device. func RemovePreset(ds *datastore.DataStore, account, device string, presetNumber int) error { - presets, err := ds.GetPresets(account, device) - if err != nil { - return err - } + _, err := ds.MutatePresets(account, device, func(presets []models.ServicePreset) ([]models.ServicePreset, error) { + if presetNumber < 1 || presetNumber > len(presets) { + // Preset doesn't exist or index out of range, nothing to do + return presets, nil + } - if presetNumber < 1 || presetNumber > len(presets) { - // Preset doesn't exist or index out of range, nothing to do - return nil - } + presets[presetNumber-1] = models.ServicePreset{} - presets[presetNumber-1] = models.ServicePreset{} + return presets, nil + }) - return ds.SavePresets(account, device, presets) + return err } // resolvePresetSource resolves the source a preset PUT is referencing, @@ -1606,11 +1605,6 @@ func UpdatePreset(ds *datastore.DataStore, account, device string, presetNumber return nil, err } - presets, err := ds.GetPresets(account, device) - if err != nil { - presets = []models.ServicePreset{} - } - var newPresetElem struct { Name string `xml:"name"` Username string `xml:"username"` @@ -1676,14 +1670,19 @@ func UpdatePreset(ds *datastore.DataStore, account, device string, presetNumber Username: newPresetElem.Name, } - // Ensure presets list is large enough - for len(presets) < presetNumber { - presets = append(presets, models.ServicePreset{}) - } + // Read-mutate-write atomically: a concurrent PUT for a different preset + // number racing this one must not be able to clobber it. See + // MutatePresets — this is the exact interleave that dropped a preset + // during #614's rapid-fire repro. + if _, err = ds.MutatePresets(account, device, func(presets []models.ServicePreset) ([]models.ServicePreset, error) { + for len(presets) < presetNumber { + presets = append(presets, models.ServicePreset{}) + } - presets[presetNumber-1] = presetObj + presets[presetNumber-1] = presetObj - if err = ds.SavePresets(account, device, presets); err != nil { + return presets, nil + }); err != nil { return nil, err } @@ -1812,11 +1811,6 @@ func AddRecent(ds *datastore.DataStore, account, device string, sourceXML []byte return nil, err } - recents, err := ds.GetRecents(account, device) - if err != nil && !os.IsNotExist(err) { - return nil, err - } - var input recentInput if err := xml.Unmarshal(sourceXML, &input); err != nil { return nil, err @@ -1861,9 +1855,19 @@ func AddRecent(ds *datastore.DataStore, account, device string, sourceXML []byte syncMatchingSource(matchingSrc, input) utcTime := parseLastPlayedAt(input.LastPlayedAt) - recentObj, recents := updateOrCreateRecent(recents, input.Name, matchingSrc, input.ContentItemType, input.Location, device, utcTime) - if err := ds.SaveRecents(account, device, recents); err != nil { + // Read-mutate-write atomically: a concurrent AddRecent/preset call for + // the same device racing this one must not be able to clobber it. See + // MutatePresets/MutateRecents for why a plain GetRecents+SaveRecents + // isn't safe here. + var recentObj *models.ServiceRecent + + if _, err := ds.MutateRecents(account, device, func(recents []models.ServiceRecent) ([]models.ServiceRecent, error) { + var updated []models.ServiceRecent + recentObj, updated = updateOrCreateRecent(recents, input.Name, matchingSrc, input.ContentItemType, input.Location, device, utcTime) + + return updated, nil + }); err != nil { return nil, err } @@ -1891,7 +1895,7 @@ func learnSource(ds *datastore.DataStore, account, device string, sources []mode matchingSrc.SecretType = constants.CredentialTypeToken } - persistLearnedSource(ds, account, device, sources, matchingSrc) + persistLearnedSource(ds, account, device, matchingSrc) } return matchingSrc, sourceLearned @@ -2041,26 +2045,24 @@ func updateSourceFields(src *models.ConfiguredSource, credentialValue, sourceNam return learned } -func persistLearnedSource(ds *datastore.DataStore, account, device string, sources []models.ConfiguredSource, matchingSrc *models.ConfiguredSource) { - updatedSources := make([]models.ConfiguredSource, len(sources)) - copy(updatedSources, sources) +func persistLearnedSource(ds *datastore.DataStore, account, device string, matchingSrc *models.ConfiguredSource) { + // Read-mutate-write atomically against the persisted list, not a + // snapshot the caller read earlier — AddRecent and UpdatePreset can + // both be learning/auto-adding sources for the same device + // concurrently, and a plain Get+Save here would silently lose + // whichever write landed second. + _, err := ds.MutateConfiguredSources(account, device, func(sources []models.ConfiguredSource) ([]models.ConfiguredSource, error) { + for i := range sources { + if sources[i].ID == matchingSrc.ID { + sources[i] = *matchingSrc - found := false - - for i := range updatedSources { - if updatedSources[i].ID == matchingSrc.ID { - updatedSources[i] = *matchingSrc - found = true - - break + return sources, nil + } } - } - if !found { - updatedSources = append(updatedSources, *matchingSrc) - } - - if err := ds.SaveConfiguredSources(account, device, updatedSources); err != nil { + return append(sources, *matchingSrc), nil + }) + if err != nil { log.Printf("[MARGE_ERR] Failed to persist learned source for %s: %s", sanitizeLog(device), sanitizeErr(err)) } } @@ -2446,7 +2448,6 @@ func AddSource(ds *datastore.DataStore, account, username, providerID, secret, s } devID := entry.Name() - sources, _ := ds.GetConfiguredSources(account, devID) newSrc := models.ConfiguredSource{ ID: sourceID, @@ -2476,37 +2477,39 @@ func AddSource(ds *datastore.DataStore, account, username, providerID, secret, s PrepareConfiguredSource(&newSrc) - // Update or append. Most providers are singletons (one account each), so - // the same provider replaces the existing entry. STORED_MUSIC is the - // exception: each DLNA media server is a separate account (username = - // "/0"), so it must only replace when the account also matches. - // Otherwise registering a second media server overwrites the first, which - // then vanishes from /full + /sources and the speaker drops it (only one - // media server could ever stay registered). - replaced := false + // Read-mutate-write atomically against the persisted list, not a + // snapshot read before the loop body — see MutateConfiguredSources. + _, err := ds.MutateConfiguredSources(account, devID, func(sources []models.ConfiguredSource) ([]models.ConfiguredSource, error) { + // Update or append. Most providers are singletons (one account + // each), so the same provider replaces the existing entry. + // STORED_MUSIC is the exception: each DLNA media server is a + // separate account (username = "/0"), so it must only + // replace when the account also matches. Otherwise registering + // a second media server overwrites the first, which then + // vanishes from /full + /sources and the speaker drops it + // (only one media server could ever stay registered). + for i := range sources { + sameProvider := sources[i].SourceProviderID == providerID + if providerID == strconv.Itoa(constants.StoredMusicProviderID) { + // Match on the persisted account identity + // (SourceKey.Account), not Username, which does not + // round-trip through the datastore. + sameProvider = sameProvider && sources[i].SourceKey.Account == username + } - for i := range sources { - sameProvider := sources[i].SourceProviderID == providerID - if providerID == strconv.Itoa(constants.StoredMusicProviderID) { - // Match on the persisted account identity (SourceKey.Account), - // not Username, which does not round-trip through the datastore. - sameProvider = sameProvider && sources[i].SourceKey.Account == username + if sameProvider || + (providerID == strconv.Itoa(constants.SpotifyProviderID) && sources[i].SourceKey.Type == constants.ProviderSpotify) { + sources[i] = newSrc + + return sources, nil + } } - if sameProvider || - (providerID == strconv.Itoa(constants.SpotifyProviderID) && sources[i].SourceKey.Type == constants.ProviderSpotify) { - sources[i] = newSrc - replaced = true - - break - } + return append(sources, newSrc), nil + }) + if err != nil { + log.Printf("[Marge] AddSource: failed to save source %s for device %s: %s", sanitizeLog(newSrc.SourceKey.Type), sanitizeLog(devID), sanitizeErr(err)) } - - if !replaced { - sources = append(sources, newSrc) - } - - _ = ds.SaveConfiguredSources(account, devID, sources) } return sourceID, nil