104 lines
2.5 KiB
Go
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()
|
|
}
|