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() }