From 2df9e53aa1796894901e0cc056b7bcc448abb933 Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Mon, 13 Apr 2026 12:26:38 -0700 Subject: [PATCH] automod: limit blob processing concurrency per record --- automod/engine/engine.go | 3 ++ automod/engine/ruleset.go | 73 ++++++++++++++------------------------- 2 files changed, 29 insertions(+), 47 deletions(-) diff --git a/automod/engine/engine.go b/automod/engine/engine.go index c6d8bab6..ff223d8e 100644 --- a/automod/engine/engine.go +++ b/automod/engine/engine.go @@ -61,6 +61,9 @@ type EngineConfig struct { IdentityEventTimeout time.Duration // timeout for event processing (total, including all setup, rules, and teardown) OzoneEventTimeout time.Duration + + // number of blobs fetched and processed concurrently per record (there may be many records processed in parallel). If zero, a sane default is used (`DEFAULT_BLOB_PROCESSING_CONCURRENCY`). If negative, no limit. + BlobProcessingConcurrency int } // Entrypoint for external code pushing #identity events in to the engine. diff --git a/automod/engine/ruleset.go b/automod/engine/ruleset.go index 1d06e23a..59dc3645 100644 --- a/automod/engine/ruleset.go +++ b/automod/engine/ruleset.go @@ -3,12 +3,15 @@ package engine import ( "bytes" "fmt" - "sync" appbsky "github.com/bluesky-social/indigo/api/bsky" lexutil "github.com/bluesky-social/indigo/lex/util" + + "golang.org/x/sync/errgroup" ) +var DEFAULT_BLOB_PROCESSING_CONCURRENCY = 8 + // Holds configuration of which rules of various types should be run, and helps dispatch events to those rules. type RuleSet struct { PostRules []PostRuleFunc @@ -121,60 +124,36 @@ func (r *RuleSet) fetchAndProcessBlobs(c *RecordContext) error { return nil } - errChan := make(chan error, len(blobs)) - var wg sync.WaitGroup + concurrency := c.engine.Config.BlobProcessingConcurrency + if concurrency == 0 { + concurrency = DEFAULT_BLOB_PROCESSING_CONCURRENCY + } + + var wg errgroup.Group + wg.SetLimit(concurrency) for _, blob := range blobs { - wg.Add(1) - go func(blob lexutil.LexBlob) { - defer wg.Done() + blob := blob + wg.Go(func() error { data, err := c.fetchBlob(blob) if err != nil { - errChan <- err - return - } - err = r.processBlob(c, blob, data) - if err != nil { - errChan <- err - return + return err } - }(blob) - - } - wg.Wait() - close(errChan) - - // check for errors - for err := range errChan { - if err != nil { - return err - } + return r.processBlob(c, blob, data) + }) } - return nil + return wg.Wait() } func (r *RuleSet) processBlob(c *RecordContext, blob lexutil.LexBlob, data []byte) error { - errChan := make(chan error, len(r.BlobRules)) - var wg sync.WaitGroup - for _, f := range r.BlobRules { - wg.Add(1) - go func(brf BlobRuleFunc) { - defer wg.Done() - err := brf(c, blob, data) - if err != nil { - errChan <- err - return - } - }(f) - } - wg.Wait() - close(errChan) - - // check for errors - for err := range errChan { - if err != nil { - return err - } + // note that this errgroup.Group is *not* bounded: it runs all blob rules in parallel + var wg errgroup.Group + for _, brf := range r.BlobRules { + brf := brf + wg.Go(func() error { + return brf(c, blob, data) + }) } - return nil + + return wg.Wait() } -- 2.51.2