forked from MirrorHub/mautrix-whatsapp
Use single SQL query for fetching backfill queue items
This commit is contained in:
parent
86840dec59
commit
6b69ea707b
4 changed files with 37 additions and 17 deletions
|
@ -20,6 +20,7 @@ import (
|
|||
"time"
|
||||
|
||||
log "maunium.net/go/maulogger/v2"
|
||||
|
||||
"maunium.net/go/mautrix-whatsapp/database"
|
||||
)
|
||||
|
||||
|
@ -32,19 +33,16 @@ type BackfillQueue struct {
|
|||
log log.Logger
|
||||
}
|
||||
|
||||
// Immediate backfills should happen first, then deferred backfills and lastly
|
||||
// media backfills.
|
||||
// RunLoop fetches backfills from the database, prioritizing immediate and forward backfills
|
||||
func (bq *BackfillQueue) RunLoop(user *User) {
|
||||
for {
|
||||
if immediate := bq.BackfillQuery.GetNext(user.MXID, database.BackfillImmediate); immediate != nil {
|
||||
bq.ImmediateBackfillRequests <- immediate
|
||||
immediate.MarkDone()
|
||||
} else if backfill := bq.BackfillQuery.GetNext(user.MXID, database.BackfillDeferred); backfill != nil {
|
||||
bq.DeferredBackfillRequests <- backfill
|
||||
if backfill := bq.BackfillQuery.GetNext(user.MXID); backfill != nil {
|
||||
if backfill.BackfillType == database.BackfillImmediate || backfill.BackfillType == database.BackfillForward {
|
||||
bq.ImmediateBackfillRequests <- backfill
|
||||
} else {
|
||||
bq.DeferredBackfillRequests <- backfill
|
||||
}
|
||||
backfill.MarkDone()
|
||||
} else if mediaBackfill := bq.BackfillQuery.GetNext(user.MXID, database.BackfillMedia); mediaBackfill != nil {
|
||||
bq.DeferredBackfillRequests <- mediaBackfill
|
||||
mediaBackfill.MarkDone()
|
||||
} else {
|
||||
select {
|
||||
case <-bq.ReCheckQueue:
|
||||
|
|
|
@ -30,14 +30,17 @@ type BackfillType int
|
|||
|
||||
const (
|
||||
BackfillImmediate BackfillType = 0
|
||||
BackfillDeferred = 1
|
||||
BackfillMedia = 2
|
||||
BackfillForward BackfillType = 100
|
||||
BackfillDeferred BackfillType = 200
|
||||
BackfillMedia BackfillType = 300
|
||||
)
|
||||
|
||||
func (bt BackfillType) String() string {
|
||||
switch bt {
|
||||
case BackfillImmediate:
|
||||
return "IMMEDIATE"
|
||||
case BackfillForward:
|
||||
return "FORWARD"
|
||||
case BackfillDeferred:
|
||||
return "DEFERRED"
|
||||
case BackfillMedia:
|
||||
|
@ -80,16 +83,15 @@ const (
|
|||
SELECT queue_id, user_mxid, type, priority, portal_jid, portal_receiver, time_start, time_end, max_batch_events, max_total_events, batch_delay
|
||||
FROM backfill_queue
|
||||
WHERE user_mxid=$1
|
||||
AND type=$2
|
||||
AND completed_at IS NULL
|
||||
ORDER BY priority, queue_id
|
||||
ORDER BY type, priority, queue_id
|
||||
LIMIT 1
|
||||
`
|
||||
)
|
||||
|
||||
// GetNext returns the next backfill to perform
|
||||
func (bq *BackfillQuery) GetNext(userID id.UserID, backfillType BackfillType) (backfill *Backfill) {
|
||||
rows, err := bq.db.Query(getNextBackfillQuery, userID, backfillType)
|
||||
func (bq *BackfillQuery) GetNext(userID id.UserID) (backfill *Backfill) {
|
||||
rows, err := bq.db.Query(getNextBackfillQuery, userID)
|
||||
defer rows.Close()
|
||||
if err != nil || rows == nil {
|
||||
bq.log.Error(err)
|
||||
|
|
20
database/upgrades/2022-04-29-backfillqueue-type-order.go
Normal file
20
database/upgrades/2022-04-29-backfillqueue-type-order.go
Normal file
|
@ -0,0 +1,20 @@
|
|||
package upgrades
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
)
|
||||
|
||||
func init() {
|
||||
upgrades[41] = upgrade{"Update backfill queue tables to be sortable by priority", func(tx *sql.Tx, ctx context) error {
|
||||
_, err := tx.Exec(`
|
||||
UPDATE backfill_queue
|
||||
SET type=CASE
|
||||
WHEN type=1 THEN 200
|
||||
WHEN type=2 THEN 300
|
||||
ELSE type
|
||||
END
|
||||
WHERE type=1 OR type=2
|
||||
`)
|
||||
return err
|
||||
}}
|
||||
}
|
|
@ -40,7 +40,7 @@ type upgrade struct {
|
|||
fn upgradeFunc
|
||||
}
|
||||
|
||||
const NumberOfUpgrades = 41
|
||||
const NumberOfUpgrades = 42
|
||||
|
||||
var upgrades [NumberOfUpgrades]upgrade
|
||||
|
||||
|
|
Loading…
Reference in a new issue