From 3055f6d317b73878a25e17db24e20e4b249d21d9 Mon Sep 17 00:00:00 2001 From: gohumble Date: Tue, 30 Aug 2022 23:00:07 +0200 Subject: [PATCH 1/4] mutex state --- internal/telegram/generate.go | 22 ++++++++++++++++++---- 1 file changed, 18 insertions(+), 4 deletions(-) diff --git a/internal/telegram/generate.go b/internal/telegram/generate.go index 2a4e0fd..ca7d054 100644 --- a/internal/telegram/generate.go +++ b/internal/telegram/generate.go @@ -98,6 +98,8 @@ func (bot *TipBot) confirmGenerateImages(ctx intercept.Context) (intercept.Conte return ctx, nil } +var mutexStates = map[int]bool{0: false, 1: false} + // generateDalleImages is called by the invoice event when the user has paid func (bot *TipBot) generateDalleImages(event Event) { invoiceEvent := event.(*InvoiceEvent) @@ -107,10 +109,22 @@ func (bot *TipBot) generateDalleImages(event Event) { return } bot.trySendMessage(user.Telegram, "🔄 Your images are being generated. Please wait a few moments.") - - // we can have only one user using dalle - mutex.Lock("dalle-image-task") - defer mutex.Unlock("dalle-image-task") + locker := -1 + for i := 0; i < 2; i++ { + if mutexStates[i] == false { + mutexStates[i] = true + locker = i + mutex.Lock(fmt.Sprintf("dalle-image-task-%d", i)) + break + } + if i == len(mutexStates)-1 { + time.Sleep(time.Second * 1) + i = 0 + } + } + if locker > 0 { + defer mutex.Unlock(fmt.Sprintf("dalle-image-task-%d", locker)) + } time.Sleep(time.Second * 1) // create the client with the bearer token api key From f288bb05f624db2c921b3161dcaf191739ab99fe Mon Sep 17 00:00:00 2001 From: gohumble Date: Tue, 30 Aug 2022 23:23:40 +0200 Subject: [PATCH 2/4] added jobchan and 2 workers --- internal/telegram/generate.go | 137 ++++++++++++++++------------------ 1 file changed, 66 insertions(+), 71 deletions(-) diff --git a/internal/telegram/generate.go b/internal/telegram/generate.go index ca7d054..2523a6b 100644 --- a/internal/telegram/generate.go +++ b/internal/telegram/generate.go @@ -13,7 +13,6 @@ import ( "github.com/LightningTipBot/LightningTipBot/internal/dalle" "github.com/LightningTipBot/LightningTipBot/internal/lnbits" "github.com/LightningTipBot/LightningTipBot/internal/runtime" - "github.com/LightningTipBot/LightningTipBot/internal/runtime/mutex" "github.com/LightningTipBot/LightningTipBot/internal/telegram/intercept" log "github.com/sirupsen/logrus" "github.com/skip2/go-qrcode" @@ -98,7 +97,18 @@ func (bot *TipBot) confirmGenerateImages(ctx intercept.Context) (intercept.Conte return ctx, nil } -var mutexStates = map[int]bool{0: false, 1: false} +var jobChan chan func() + +func init() { + jobChan = make(chan func(), 2) + go worker(jobChan) + go worker(jobChan) +} +func worker(linkChan chan func()) { + for generatePrompt := range linkChan { + generatePrompt() + } +} // generateDalleImages is called by the invoice event when the user has paid func (bot *TipBot) generateDalleImages(event Event) { @@ -109,86 +119,71 @@ func (bot *TipBot) generateDalleImages(event Event) { return } bot.trySendMessage(user.Telegram, "🔄 Your images are being generated. Please wait a few moments.") - locker := -1 - for i := 0; i < 2; i++ { - if mutexStates[i] == false { - mutexStates[i] = true - locker = i - mutex.Lock(fmt.Sprintf("dalle-image-task-%d", i)) - break - } - if i == len(mutexStates)-1 { - time.Sleep(time.Second * 1) - i = 0 - } - } - if locker > 0 { - defer mutex.Unlock(fmt.Sprintf("dalle-image-task-%d", locker)) - } - time.Sleep(time.Second * 1) - - // create the client with the bearer token api key - dalleClient, err := dalle.NewHTTPClient(internal.Configuration.Generate.DalleKey) - // handle err - if err != nil { - log.Errorf("[NewHTTPClient] %v", err.Error()) - bot.dalleRefundUser(user) - return - } - - ctx, cancel := context.WithTimeout(context.Background(), time.Minute*10) - defer cancel() - // generate a task to create an image with a prompt - task, err := dalleClient.Generate(ctx, invoiceEvent.CallbackData) - if err != nil { - log.Errorf("[Generate] %v", err.Error()) - bot.dalleRefundUser(user) - return - } - // poll the task.ID until status is succeeded - var t *dalle.Task - timeout := time.After(5 * time.Minute) - ticker := time.Tick(5 * time.Second) - // Keep trying until we're timed out or get a result/error - for { - select { - case <-ctx.Done(): + var job = func() { + // create the client with the bearer token api key + dalleClient, err := dalle.NewHTTPClient(internal.Configuration.Generate.DalleKey) + // handle err + if err != nil { + log.Errorf("[NewHTTPClient] %v", err.Error()) bot.dalleRefundUser(user) - log.Errorf("[DALLE] ctx done") return - // Got a timeout! fail with a timeout error - case <-timeout: + } + + ctx, cancel := context.WithTimeout(context.Background(), time.Minute*10) + defer cancel() + // generate a task to create an image with a prompt + task, err := dalleClient.Generate(ctx, invoiceEvent.CallbackData) + if err != nil { + log.Errorf("[Generate] %v", err.Error()) bot.dalleRefundUser(user) - log.Errorf("[DALLE] timeout") return - // Got a tick, we should check on checkSomething() - case <-ticker: - t, err = dalleClient.GetTask(ctx, task.ID) - // handle err - if err != nil { - log.Errorf("[GetTask] %v", err.Error()) + } + // poll the task.ID until status is succeeded + var t *dalle.Task + timeout := time.After(5 * time.Minute) + ticker := time.Tick(5 * time.Second) + // Keep trying until we're timed out or get a result/error + for { + select { + case <-ctx.Done(): bot.dalleRefundUser(user) + log.Errorf("[DALLE] ctx done") return - } - if t.Status == dalle.StatusSucceeded { - fmt.Printf("[DALLE] task succeeded for user %s", GetUserStr(user.Telegram)) - // download the first generated image - for _, data := range t.Generations.Data { - err = bot.downloadAndSendImages(ctx, dalleClient, data, invoiceEvent) - if err != nil { - log.Errorf("[downloadAndSendImages] %v", err.Error()) - } + // Got a timeout! fail with a timeout error + case <-timeout: + bot.dalleRefundUser(user) + log.Errorf("[DALLE] timeout") + return + // Got a tick, we should check on checkSomething() + case <-ticker: + t, err = dalleClient.GetTask(ctx, task.ID) + // handle err + if err != nil { + log.Errorf("[GetTask] %v", err.Error()) + bot.dalleRefundUser(user) + return } - return + if t.Status == dalle.StatusSucceeded { + fmt.Printf("[DALLE] task succeeded for user %s", GetUserStr(user.Telegram)) + // download the first generated image + for _, data := range t.Generations.Data { + err = bot.downloadAndSendImages(ctx, dalleClient, data, invoiceEvent) + if err != nil { + log.Errorf("[downloadAndSendImages] %v", err.Error()) + } + } + return - } else if t.Status == dalle.StatusRejected { - log.Errorf("[DALLE] rejected: %s", t.ID) - bot.dalleRefundUser(user) - return + } else if t.Status == dalle.StatusRejected { + log.Errorf("[DALLE] rejected: %s", t.ID) + bot.dalleRefundUser(user) + return + } + log.Debugf("[DALLE] pending for user %s", GetUserStr(user.Telegram)) } - log.Debugf("[DALLE] pending for user %s", GetUserStr(user.Telegram)) } } + jobChan <- job } // downloadAndSendImages will download dalle images and send them to the payer. From 89552a73b7642f90da93c8f61bc7f2e2a2c0665e Mon Sep 17 00:00:00 2001 From: callebtc <93376500+callebtc@users.noreply.github.com> Date: Wed, 31 Aug 2022 10:58:50 +0200 Subject: [PATCH 3/4] changes from main --- internal/telegram/generate.go | 25 ++++++++++++++++--------- 1 file changed, 16 insertions(+), 9 deletions(-) diff --git a/internal/telegram/generate.go b/internal/telegram/generate.go index 2523a6b..68719f6 100644 --- a/internal/telegram/generate.go +++ b/internal/telegram/generate.go @@ -125,7 +125,7 @@ func (bot *TipBot) generateDalleImages(event Event) { // handle err if err != nil { log.Errorf("[NewHTTPClient] %v", err.Error()) - bot.dalleRefundUser(user) + bot.dalleRefundUser(user, "") return } @@ -135,7 +135,7 @@ func (bot *TipBot) generateDalleImages(event Event) { task, err := dalleClient.Generate(ctx, invoiceEvent.CallbackData) if err != nil { log.Errorf("[Generate] %v", err.Error()) - bot.dalleRefundUser(user) + bot.dalleRefundUser(user, "") return } // poll the task.ID until status is succeeded @@ -146,12 +146,12 @@ func (bot *TipBot) generateDalleImages(event Event) { for { select { case <-ctx.Done(): - bot.dalleRefundUser(user) - log.Errorf("[DALLE] ctx done") + bot.dalleRefundUser(user, "") + log.Errorf("[DALLE] ctx done", "") return // Got a timeout! fail with a timeout error case <-timeout: - bot.dalleRefundUser(user) + bot.dalleRefundUser(user, "Timeout. Please try again later.") log.Errorf("[DALLE] timeout") return // Got a tick, we should check on checkSomething() @@ -160,7 +160,7 @@ func (bot *TipBot) generateDalleImages(event Event) { // handle err if err != nil { log.Errorf("[GetTask] %v", err.Error()) - bot.dalleRefundUser(user) + bot.dalleRefundUser(user, "") return } if t.Status == dalle.StatusSucceeded { @@ -176,7 +176,7 @@ func (bot *TipBot) generateDalleImages(event Event) { } else if t.Status == dalle.StatusRejected { log.Errorf("[DALLE] rejected: %s", t.ID) - bot.dalleRefundUser(user) + bot.dalleRefundUser(user, "Your prompt has been rejected by OpenAI. Do not use celebrity names, sexual expressions, or any other harmful content as prompt.") return } log.Debugf("[DALLE] pending for user %s", GetUserStr(user.Telegram)) @@ -212,7 +212,7 @@ func (bot *TipBot) downloadAndSendImages(ctx context.Context, dalleClient dalle. return nil } -func (bot *TipBot) dalleRefundUser(user *lnbits.User) error { +func (bot *TipBot) dalleRefundUser(user *lnbits.User, message string) error { if user.Wallet == nil { return fmt.Errorf("user has no wallet") } @@ -240,6 +240,13 @@ func (bot *TipBot) dalleRefundUser(user *lnbits.User) error { return err } log.Warnf("[DALLE] refunding user %s with %d sat", GetUserStr(user.Telegram), internal.Configuration.Generate.DallePrice) - bot.trySendMessage(user.Telegram, "🚫 Something went wrong. You have been refunded.") + + var err_reason string + if len(message) > 0 { + err_reason = message + } else { + err_reason = "Something went wrong." + } + bot.trySendMessage(user.Telegram, fmt.Sprintf("🚫 %s You have been refunded.", err_reason)) return nil } From e170394eb8c6fdf1a67f2aff791184a98053792a Mon Sep 17 00:00:00 2001 From: gohumble Date: Wed, 31 Aug 2022 15:44:29 +0200 Subject: [PATCH 4/4] add workerId --- internal/telegram/generate.go | 32 +++++++++++++++++--------------- 1 file changed, 17 insertions(+), 15 deletions(-) diff --git a/internal/telegram/generate.go b/internal/telegram/generate.go index 2523a6b..4893cdc 100644 --- a/internal/telegram/generate.go +++ b/internal/telegram/generate.go @@ -97,16 +97,18 @@ func (bot *TipBot) confirmGenerateImages(ctx intercept.Context) (intercept.Conte return ctx, nil } -var jobChan chan func() +var jobChan chan func(workerId int) +var workers = 2 func init() { - jobChan = make(chan func(), 2) - go worker(jobChan) - go worker(jobChan) + jobChan = make(chan func(workerId int), workers) + for i := 0; i < workers; i++ { + go worker(jobChan, i) + } } -func worker(linkChan chan func()) { +func worker(linkChan chan func(workerId int), workerId int) { for generatePrompt := range linkChan { - generatePrompt() + generatePrompt(workerId) } } @@ -119,12 +121,12 @@ func (bot *TipBot) generateDalleImages(event Event) { return } bot.trySendMessage(user.Telegram, "🔄 Your images are being generated. Please wait a few moments.") - var job = func() { + var job = func(worker int) { // create the client with the bearer token api key dalleClient, err := dalle.NewHTTPClient(internal.Configuration.Generate.DalleKey) // handle err if err != nil { - log.Errorf("[NewHTTPClient] %v", err.Error()) + log.Errorf("[NewHTTPClient-%d] %v", worker, err.Error()) bot.dalleRefundUser(user) return } @@ -134,7 +136,7 @@ func (bot *TipBot) generateDalleImages(event Event) { // generate a task to create an image with a prompt task, err := dalleClient.Generate(ctx, invoiceEvent.CallbackData) if err != nil { - log.Errorf("[Generate] %v", err.Error()) + log.Errorf("[Generate-%d] %v", worker, err.Error()) bot.dalleRefundUser(user) return } @@ -147,12 +149,12 @@ func (bot *TipBot) generateDalleImages(event Event) { select { case <-ctx.Done(): bot.dalleRefundUser(user) - log.Errorf("[DALLE] ctx done") + log.Errorf("[DALLE-%d] ctx done", worker) return // Got a timeout! fail with a timeout error case <-timeout: bot.dalleRefundUser(user) - log.Errorf("[DALLE] timeout") + log.Errorf("[DALLE-%d] timeout", worker) return // Got a tick, we should check on checkSomething() case <-ticker: @@ -164,22 +166,22 @@ func (bot *TipBot) generateDalleImages(event Event) { return } if t.Status == dalle.StatusSucceeded { - fmt.Printf("[DALLE] task succeeded for user %s", GetUserStr(user.Telegram)) + fmt.Printf("[DALLE-%d] task succeeded for user %s", worker, GetUserStr(user.Telegram)) // download the first generated image for _, data := range t.Generations.Data { err = bot.downloadAndSendImages(ctx, dalleClient, data, invoiceEvent) if err != nil { - log.Errorf("[downloadAndSendImages] %v", err.Error()) + log.Errorf("[downloadAndSendImages-%d] %v", worker, err.Error()) } } return } else if t.Status == dalle.StatusRejected { - log.Errorf("[DALLE] rejected: %s", t.ID) + log.Errorf("[DALLE-%d] rejected: %s", worker, t.ID) bot.dalleRefundUser(user) return } - log.Debugf("[DALLE] pending for user %s", GetUserStr(user.Telegram)) + log.Debugf("[DALLE-%d] pending for user %s", worker, GetUserStr(user.Telegram)) } } }