@@ -15,6 +15,7 @@ import (
1515 "context"
1616 "fmt"
1717 "sort"
18+ "strings"
1819 "sync"
1920
2021 "github.com/ipfs/go-cid"
@@ -76,6 +77,12 @@ func (db *DB) Merge(ctx context.Context, evt event.Merge) error {
7677// several, so a transaction's cost grows with the square of the events in it.
7778const mergeChunkSize = 8
7879
80+ // Phases a merge chunk can fail in, reported on the retry-exhaustion log line.
81+ const (
82+ phaseRead = "read"
83+ phaseCommit = "commit"
84+ )
85+
7986type mergeEntry struct {
8087 evt event.Merge
8188 col * collection
@@ -106,6 +113,7 @@ func (db *DB) MergeBatchWithTxn(ctx context.Context, merges []event.Merge) ([]bo
106113 col , err := getCollectionFromCollectionID (ctx , db , evt .CollectionID )
107114 if err != nil {
108115 errs = append (errs , NewErrMergeEventDropped (err , evt .DocID , evt .Cid .String ()))
116+ db .stats .markDropped (dropCollection )
109117 continue
110118 }
111119 entries = append (entries , mergeEntry {evt : evt , col : col , index : i })
@@ -178,6 +186,7 @@ func (db *DB) MergeBatchWithTxn(ctx context.Context, merges []event.Merge) ([]bo
178186 for i := range chunk {
179187 if err := db .mergeChunk (ctx , chunk [i :i + 1 ]); err != nil {
180188 errs = append (errs , NewErrMergeEventDropped (err , chunk [i ].evt .DocID , chunk [i ].evt .Cid .String ()))
189+ db .stats .markDropped (mergeDropReason (err ))
181190 continue
182191 }
183192 db .publishMergeComplete (chunk [i : i + 1 ])
@@ -197,29 +206,50 @@ func (db *DB) txnAttempts() int {
197206 return 1
198207}
199208
209+ // namedDocs renders the documents a transaction touched, as collection/docID. Badger
210+ // reports conflicts without naming the contended key, so this is the only lead available
211+ // for working out which documents contend with each other.
212+ func namedDocs (entries []mergeEntry ) string {
213+ docIDs := make ([]string , len (entries ))
214+ for i , e := range entries {
215+ docIDs [i ] = e .col .Name () + "/" + e .evt .DocID
216+ }
217+ return strings .Join (docIDs , "," )
218+ }
219+
200220// mergeChunk merges every event of the chunk inside one transaction, retrying the
201221// whole chunk on transaction conflict. Isolating a failing event is the caller's job.
202222func (db * DB ) mergeChunk (ctx context.Context , entries []mergeEntry ) error {
203- // Held so that exhausting the retry budget can report the conflict that caused it.
223+ // Held so that exhausting the retry budget can report the conflict that caused it and
224+ // where the last one was raised.
204225 var conflictErr error
226+ var phase string
227+ // Whether each event created its document, kept until the transaction commits so a
228+ // retried attempt does not count its events twice.
229+ creates := make ([]bool , 0 , len (entries ))
205230 for i := 0 ; i < db .txnAttempts (); i ++ {
206231 txn , err := db .NewTxn (false )
207232 if err != nil {
208233 return err
209234 }
210235 txnCtx := InitContext (ctx , txn )
211236
237+ creates = creates [:0 ]
212238 var mergeErr error
213239 for _ , e := range entries {
214- if mergeErr = db .mergeInTxn (txnCtx , e .col , e .evt ); mergeErr != nil {
240+ isCreate , err := db .mergeInTxn (txnCtx , e .col , e .evt )
241+ if err != nil {
242+ mergeErr , phase = err , phaseRead
215243 break
216244 }
245+ creates = append (creates , isCreate )
217246 }
218247
219248 if mergeErr != nil {
220249 txn .Discard ()
221250 if errors .Is (mergeErr , corekv .ErrTxnConflict ) {
222251 conflictErr = mergeErr
252+ db .stats .chunkConflicts .Add (1 )
223253 continue
224254 }
225255 return mergeErr
@@ -229,15 +259,29 @@ func (db *DB) mergeChunk(ctx context.Context, entries []mergeEntry) error {
229259 txn .Discard ()
230260 if errors .Is (err , corekv .ErrTxnConflict ) {
231261 conflictErr = err
262+ db .stats .chunkConflicts .Add (1 )
263+ phase = phaseCommit
232264 continue
233265 }
234266 return err
235267 }
236268
269+ for _ , isCreate := range creates {
270+ db .stats .markCreateOrUpdate (isCreate )
271+ }
237272 return nil
238273 }
239274
240- // Nothing was committed, so callers must not treat the events as merged.
275+ // The chunk used its whole retry budget without committing. The caller then re-runs it
276+ // one event at a time and most events usually land on that pass, so this counts
277+ // conflict pressure, not loss. What was actually lost is named in the caller's error.
278+ db .stats .markExhausted ()
279+
280+ log .InfoContext (ctx , "merge chunk exhausted its retries" ,
281+ corelog .Int ("attempts" , db .txnAttempts ()),
282+ corelog .String ("phase" , phase ),
283+ corelog .String ("docIDs" , namedDocs (entries )),
284+ )
241285 return client .NewErrMaxTxnRetries (conflictErr )
242286}
243287
@@ -254,13 +298,15 @@ func (db *DB) executeMerge(ctx context.Context, col *collection, dagMerge event.
254298 }
255299 defer txn .Discard ()
256300
257- if err := db .mergeInTxn (ctx , col , dagMerge ); err != nil {
301+ isCreate , err := db .mergeInTxn (ctx , col , dagMerge )
302+ if err != nil {
258303 return err
259304 }
260305
261306 if err := txn .Commit (); err != nil {
262307 return err
263308 }
309+ db .stats .markCreateOrUpdate (isCreate )
264310
265311 // send a complete event so we can track merges in the integration tests
266312 db .events .Publish (event .NewMessage (event .MergeCompleteName , event.MergeComplete {Merge : dagMerge }))
@@ -269,40 +315,47 @@ func (db *DB) executeMerge(ctx context.Context, col *collection, dagMerge event.
269315
270316// mergeInTxn executes the merge logic for a single event using the transaction already
271317// present on ctx. It does not commit; the caller is responsible for committing.
272- func (db * DB ) mergeInTxn (ctx context.Context , col * collection , dagMerge event.Merge ) error {
318+ //
319+ // Reports whether the event created the document rather than updating one already held,
320+ // which the caller counts once the transaction commits.
321+ func (db * DB ) mergeInTxn (ctx context.Context , col * collection , dagMerge event.Merge ) (bool , error ) {
273322 key , exists , err := getDocHeadstoreKey (ctx , col , dagMerge .DocID )
274323 if err != nil {
275- return err
324+ return false , err
276325 }
277326
278327 mt := newMergeTarget ()
279328 if exists {
280329 mt , err = getHeadsAsMergeTarget (ctx , key )
281330 if err != nil {
282- return NewErrGetMergeTargetHeads (err , dagMerge .DocID , string (key .Bytes ()))
331+ return false , NewErrGetMergeTargetHeads (err , dagMerge .DocID , string (key .Bytes ()))
283332 }
284333 }
285334
286- mp , err := db .newMergeProcessor (ctx , col , len (mt .heads ) == 0 )
335+ // No local heads means the merge is creating the document rather than updating one
336+ // that already exists here.
337+ newDocCreateMode := len (mt .heads ) == 0
338+
339+ mp , err := db .newMergeProcessor (ctx , col , newDocCreateMode )
287340 if err != nil {
288- return err
341+ return false , err
289342 }
290343
291344 if err = mp .loadComposites (ctx , dagMerge .Cid , mt ); err != nil {
292- return NewErrLoadComposites (err , dagMerge .Cid .String (), dagMerge .DocID )
345+ return false , NewErrLoadComposites (err , dagMerge .Cid .String (), dagMerge .DocID )
293346 }
294347
295348 if err = mp .mergeComposites (ctx ); err != nil {
296- return NewErrMergeComposites (err , dagMerge .DocID )
349+ return false , NewErrMergeComposites (err , dagMerge .DocID )
297350 }
298351
299352 for docID , oldDoc := range mp .docIDs {
300353 if err = syncIndexedDoc (ctx , docID , mp .col , oldDoc ); err != nil {
301- return NewErrSyncIndexedDoc (err , docID .String ())
354+ return false , NewErrSyncIndexedDoc (err , docID .String ())
302355 }
303356 }
304357
305- return nil
358+ return newDocCreateMode , nil
306359}
307360
308361const maxConcurrentMerges = 32
0 commit comments