Files
summercms/modules/surf/serve.go
Jakub Zych b319e7cc61 feat(11-01): run job workers in serve and queue:work, add queue:clear
- Manager gains the apparatus JobManager surface: StartJob, UpdateJobState,
  UpdateMetadata, FailJob, CancelJob (is_canceled + STOPPED + River JobCancel),
  StopJob (STOPPED only), CheckIfCanceled and GetMetadata, all raw column
  writes so updated_at is untouched
- serve starts the in-process worker unless queue.work_in_serve is false and
  stops it on shutdown; an app without jobs gets an idle worker
- queue:work runs a foreground worker with repeatable --queue filters;
  queue:clear deletes available, scheduled and retryable jobs of one queue
- the generated main appends conga.RuntimeCommands; summer delegates
  queue:work and queue:clear; make:job scaffolds a conga.Job
2026-09-29 15:20:14 +02:00

124 lines
3.2 KiB
Go

package surf
import (
"context"
"fmt"
"net"
"net/http"
"os"
"os/signal"
"strings"
"syscall"
"time"
"git.golem15.com/golem15/summercms/modules/backpack"
"git.golem15.com/golem15/summercms/modules/bonfire"
"git.golem15.com/golem15/summercms/modules/conga"
"git.golem15.com/golem15/summercms/modules/lagoon"
"git.golem15.com/golem15/summercms/modules/lagoon/attach"
"git.golem15.com/golem15/summercms/modules/party"
"gocloud.dev/blob"
)
// ServeCommand starts a signal-aware HTTP server on the assembled router.
func ServeCommand(app *backpack.App, plugins []party.Plugin) bonfire.Command {
return bonfire.Command{
Name: "serve",
Description: "Serve the application HTTP API",
Flags: []bonfire.Flag{{
Name: "addr",
Description: "Listen address",
Default: ":8080",
}},
Run: func(ctx context.Context, in bonfire.Input, out bonfire.Output) error {
addr, _ := in.Flag("addr")
if strings.TrimSpace(addr) == "" {
addr = ":8080"
}
ctx, stop := signal.NotifyContext(ctx, os.Interrupt, syscall.SIGTERM)
defer stop()
sqlDB, gdb, err := lagoon.OpenFromApp(ctx, app)
if err != nil {
return err
}
defer sqlDB.Close()
if err := lagoon.Publish(app, sqlDB, gdb); err != nil {
return err
}
bucket, err := publishUploads(ctx, app)
if err != nil {
return err
}
defer bucket.Close()
h, err := Assemble(app, plugins)
if err != nil {
return err
}
worker, err := conga.StartServeWorker(ctx, app, plugins)
if err != nil {
return err
}
if worker != nil {
out.Info(fmt.Sprintf("job worker started on queues: %s", strings.Join(worker.Queues(), ", ")))
}
srv := &http.Server{Addr: addr, Handler: h, BaseContext: func(net.Listener) context.Context { return ctx }}
errCh := make(chan error, 1)
go func() {
out.Info(fmt.Sprintf("listening on %s", addr))
errCh <- srv.ListenAndServe()
}()
shutdownCtx := func() (context.Context, context.CancelFunc) {
return context.WithTimeout(context.Background(), 10*time.Second)
}
select {
case <-ctx.Done():
sctx, cancel := shutdownCtx()
defer cancel()
shutdownErr := srv.Shutdown(sctx)
workerErr := stopWorker(sctx, worker)
if shutdownErr != nil {
return shutdownErr
}
err := <-errCh
if err != nil && err != http.ErrServerClosed {
return err
}
return workerErr
case err := <-errCh:
sctx, cancel := shutdownCtx()
defer cancel()
workerErr := stopWorker(sctx, worker)
if err != nil && err != http.ErrServerClosed {
return err
}
return workerErr
}
},
}
}
// stopWorker stops the in-process job worker, if serve started one.
func stopWorker(ctx context.Context, w *conga.Worker) error {
if w == nil {
return nil
}
return w.Stop(ctx)
}
// publishUploads opens storage.uploads.bucket_url and publishes *blob.Bucket.
// An empty URL fails boot the same way an empty JWT secret does.
func publishUploads(ctx context.Context, app *backpack.App) (*blob.Bucket, error) {
if app == nil {
return nil, fmt.Errorf("surf: app is nil")
}
bucket, err := attach.OpenBucket(ctx, app.Config)
if err != nil {
return nil, err
}
if err := attach.Publish(app, bucket); err != nil {
_ = bucket.Close()
return nil, err
}
return bucket, nil
}