Files
pinger/.vendor/with/dbos/worker.go
T
2026-06-27 16:44:47 -06:00

104 lines
2.5 KiB
Go

package dbos
import (
"context"
"log"
"os"
"slices"
"time"
"github.com/dbos-inc/dbos-transact-golang/dbos"
)
type Worker struct {
dbos dbos.DBOSContext
q string
qs []dbos.WorkflowQueue
}
func QueueWorker(ctx context.Context, conn, q string, foo func(*Worker) error) error {
return NewWorker(ctx, conn, func(w *Worker) error {
return foo(w.WithQueue(q))
})
}
func NewWorker(ctx context.Context, conn string, foo func(*Worker) error) error {
type result struct {
ctx dbos.DBOSContext
err error
}
ch := make(chan result)
go func() {
defer close(ch)
dbosctx, err := dbos.NewDBOSContext(ctx, dbos.Config{
AppName: "with",
DatabaseURL: conn,
AdminServer: os.Getenv("WITH_DBOS_ADMIN_SERVER") == "true",
ApplicationVersion: "latest",
})
select {
case ch <- result{ctx: dbosctx, err: err}:
case <-ctx.Done():
}
}()
select {
case result := <-ch:
if err := result.err; err != nil {
return err
}
return foo(&Worker{dbos: result.ctx})
case <-ctx.Done():
}
return ctx.Err()
}
func (c *Worker) Listen(ctx context.Context) error {
dbos := c.dbos
//dbos.ListenQueues(c.dbos, c.qs...)
log.Printf("[worker] listening to %+v", c.q)
if err := dbos.Launch(); err != nil {
return err
}
defer dbos.Shutdown(5 * time.Second)
<-ctx.Done()
return nil
}
func (c *Worker) WithQueue(q string) *Worker {
c2 := *c
c2.q = q
c2.qs = slices.Clone(c2.qs)
dbosq := dbos.NewWorkflowQueue(c2.dbos, q)
//dbos.WithWorkerConcurrency(5), // per worker
//dbos.WithGlobalConcurrency(10), // per all workers
//dbos.WithRateLimiter(&dbos.RateLimiter{Limit: 100, Period: time.Second}),
//dbos.WithPriorityEnabled(),
// dbos.WithPartitionQueue(), // WithQueuePartitionKey but has limit per-partiiton
c2.qs = append(c2.qs, dbosq)
return &c2
}
func Can[P any, R any](c *Worker, foo dbos.Workflow[P, R]) {
log.Printf("[worker] registering %s", fooToName(foo))
dbos.RegisterWorkflow(c.dbos, foo, dbos.WithWorkflowName(fooToName(foo)))
}
func Every[P any, R any](w *Worker, foo dbos.Workflow[P, R], cron string) {
log.Printf("[worker] registering cron '%s' %s", cron, fooToName(foo))
dbos.RegisterWorkflow(w.dbos, foo, dbos.WithSchedule(cron), dbos.WithWorkflowName(fooToName(foo)))
}
func Do[P any, R any](ctx context.Context, w *Worker, foo dbos.Workflow[P, R], input P) (R, error) {
dbosctx := dbos.From(w.dbos, ctx)
log.Printf("[worker] doing %s", fooToName(foo))
handle, err := dbos.RunWorkflow(dbosctx, foo, input) //, options...)
var some R
if err != nil {
return some, err
}
return handle.GetResult()
}