Merge pull request #7694 from penpot/niwinz-staging-runner-fixes

🐛 Fix precision issues on worker task scheduling mechanism
This commit is contained in:
Alejandro Alonso 2025-11-05 12:18:23 +01:00 committed by GitHub
commit 02a1992a0a
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
4 changed files with 41 additions and 31 deletions

View File

@ -218,6 +218,9 @@
(when (or (nil? revn) (= revn (:revn file))) (when (or (nil? revn) (= revn (:revn file)))
file))) file)))
;; FIXME: we should skip files that does not match the revn on the
;; props and add proper schema for this task props
(defn- process-file! (defn- process-file!
[cfg {:keys [file-id] :as props}] [cfg {:keys [file-id] :as props}]
(if-let [file (get-file cfg props)] (if-let [file (get-file cfg props)]

View File

@ -8,6 +8,7 @@
"A maintenance task that is responsible of properly scheduling the "A maintenance task that is responsible of properly scheduling the
file-gc task for all files that matches the eligibility threshold." file-gc task for all files that matches the eligibility threshold."
(:require (:require
[app.common.logging :as l]
[app.common.time :as ct] [app.common.time :as ct]
[app.config :as cf] [app.config :as cf]
[app.db :as db] [app.db :as db]
@ -21,25 +22,24 @@
f.modified_at f.modified_at
FROM file AS f FROM file AS f
WHERE f.has_media_trimmed IS false WHERE f.has_media_trimmed IS false
AND f.modified_at < now() - ?::interval AND f.modified_at < ?
AND f.deleted_at IS NULL AND f.deleted_at IS NULL
ORDER BY f.modified_at DESC ORDER BY f.modified_at DESC
FOR UPDATE OF f FOR UPDATE OF f
SKIP LOCKED") SKIP LOCKED")
(defn- get-candidates
[{:keys [::db/conn ::min-age] :as cfg}]
(let [min-age (db/interval min-age)]
(db/plan conn [sql:get-candidates min-age] {:fetch-size 10})))
(defn- schedule! (defn- schedule!
[cfg] [{:keys [::db/conn] :as cfg} threshold]
(let [total (reduce (fn [total {:keys [id modified-at revn]}] (let [total (reduce (fn [total {:keys [id modified-at revn]}]
(let [params {:file-id id :modified-at modified-at :revn revn}] (let [params {:file-id id :revn revn}]
(l/trc :hint "schedule"
:file-id (str id)
:revn revn
:modified-at (ct/format-inst modified-at))
(wrk/submit! (assoc cfg ::wrk/params params)) (wrk/submit! (assoc cfg ::wrk/params params))
(inc total))) (inc total)))
0 0
(get-candidates cfg))] (db/plan conn [sql:get-candidates threshold] {:fetch-size 10}))]
{:processed total})) {:processed total}))
(defmethod ig/assert-key ::handler (defmethod ig/assert-key ::handler
@ -53,12 +53,12 @@
(defmethod ig/init-key ::handler (defmethod ig/init-key ::handler
[_ cfg] [_ cfg]
(fn [{:keys [props] :as task}] (fn [{:keys [props] :as task}]
(let [min-age (ct/duration (or (:min-age props) (::min-age cfg)))] (let [threshold (-> (ct/duration (or (:min-age props) (::min-age cfg)))
(ct/in-past))]
(-> cfg (-> cfg
(assoc ::db/rollback (:rollback? props)) (assoc ::db/rollback (:rollback? props))
(assoc ::min-age min-age)
(assoc ::wrk/task :file-gc) (assoc ::wrk/task :file-gc)
(assoc ::wrk/priority 10) (assoc ::wrk/priority 10)
(assoc ::wrk/mark-retries 0) (assoc ::wrk/mark-retries 0)
(assoc ::wrk/delay 1000) (assoc ::wrk/delay 10000)
(db/tx-run! schedule!))))) (db/tx-run! schedule! threshold)))))

View File

@ -77,8 +77,8 @@
;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;; ;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;
(def ^:private sql:insert-new-task (def ^:private sql:insert-new-task
"insert into task (id, name, props, queue, label, priority, max_retries, scheduled_at) "insert into task (id, name, props, queue, label, priority, max_retries, created_at, modified_at, scheduled_at)
values (?, ?, ?, ?, ?, ?, ?, now() + ?) values (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
returning id") returning id")
(def ^:private (def ^:private
@ -88,7 +88,7 @@
AND queue=? AND queue=?
AND label=? AND label=?
AND status = 'new' AND status = 'new'
AND scheduled_at > now()") AND scheduled_at > ?")
(def ^:private schema:options (def ^:private schema:options
[:map {:title "submit-options"} [:map {:title "submit-options"}
@ -111,17 +111,19 @@
(check-options! options) (check-options! options)
(let [duration (ct/duration delay) (let [delay (ct/duration delay)
interval (db/interval duration) now (ct/now)
props (db/tjson params) scheduled-at (-> (ct/plus now delay)
id (uuid/next) (ct/truncate :millisecond))
tenant (cf/get :tenant) props (db/tjson params)
task (d/name task) id (uuid/next)
queue (str/ffmt "%:%" tenant (d/name queue)) tenant (cf/get :tenant)
conn (db/get-connectable options) task (d/name task)
deleted (when dedupe queue (str/ffmt "%:%" tenant (d/name queue))
(-> (db/exec-one! conn [sql:remove-not-started-tasks task queue label]) conn (db/get-connectable options)
:next.jdbc/update-count))] deleted (when dedupe
(-> (db/exec-one! conn [sql:remove-not-started-tasks task queue label now])
(db/get-update-count)))]
(l/trc :hint "submit task" (l/trc :hint "submit task"
:name task :name task
@ -129,11 +131,13 @@
:queue queue :queue queue
:label label :label label
:dedupe (boolean dedupe) :dedupe (boolean dedupe)
:delay (ct/format-duration duration) :delay (ct/format-duration delay)
:replace (or deleted 0)) :replace (or deleted 0))
(db/exec-one! conn [sql:insert-new-task id task props queue (db/exec-one! conn [sql:insert-new-task id task props queue
label priority max-retries interval]) label priority max-retries
now now scheduled-at])
id)) id))
(defn invoke! (defn invoke!

View File

@ -158,7 +158,9 @@
(inst-ms (:scheduled-at task))) (inst-ms (:scheduled-at task)))
(l/wrn :hint "skiping task, rescheduled" (l/wrn :hint "skiping task, rescheduled"
:task-id task-id :task-id task-id
:runner-id id) :runner-id id
:scheduled-at (ct/format-inst (:scheduled-at task))
:expected-scheduled-at (ct/format-inst scheduled-at))
:else :else
(let [result (run-task cfg task)] (let [result (run-task cfg task)]
@ -179,7 +181,8 @@
{:error explain {:error explain
:status "retry" :status "retry"
:modified-at now :modified-at now
:scheduled-at (ct/plus now delay) :scheduled-at (-> (ct/plus now delay)
(ct/truncate :millisecond))
:retry-num nretry} :retry-num nretry}
{:id (:id task)}) {:id (:id task)})
nil)) nil))