Skip to content

Commit 932065b

Browse files
committed
add periodic job endpoint
1 parent 61a03f8 commit 932065b

6 files changed

Lines changed: 150 additions & 1 deletion

File tree

riverproui/endpoints.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -175,6 +175,7 @@ func (e *endpoints[TTx]) MountEndpoints(archetype *baseservice.Archetype, logger
175175

176176
endpoints := e.ossEndpoints.MountEndpoints(archetype, logger, mux, mountOpts)
177177
endpoints = append(endpoints,
178+
apiendpoint.Mount(mux, prohandler.NewPeriodicJobListEndpoint(bundle), mountOpts),
178179
apiendpoint.Mount(mux, prohandler.NewProducerListEndpoint(bundle), mountOpts),
179180
apiendpoint.Mount(mux, prohandler.NewWorkflowCancelEndpoint(bundle), mountOpts),
180181
apiendpoint.Mount(mux, prohandler.NewWorkflowGetEndpoint(bundle), mountOpts),

riverproui/internal/prohandler/pro_handler_api_endpoints.go

Lines changed: 53 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,59 @@ func listResponseFrom[T any](data []*T) *listResponse[T] {
3939
return &listResponse[T]{Data: data}
4040
}
4141

42+
//
43+
// periodicJobListEndpoint
44+
//
45+
46+
type periodicJobListEndpoint[TTx any] struct {
47+
ProAPIBundle[TTx]
48+
apiendpoint.Endpoint[periodicJobListRequest, listResponse[uitype.RiverPeriodicJob]]
49+
}
50+
51+
func NewPeriodicJobListEndpoint[TTx any](apiBundle ProAPIBundle[TTx]) *periodicJobListEndpoint[TTx] {
52+
return &periodicJobListEndpoint[TTx]{ProAPIBundle: apiBundle}
53+
}
54+
55+
func (*periodicJobListEndpoint[TTx]) Meta() *apiendpoint.EndpointMeta {
56+
return &apiendpoint.EndpointMeta{
57+
Pattern: "GET /api/pro/periodic-jobs",
58+
StatusCode: http.StatusOK,
59+
}
60+
}
61+
62+
type periodicJobListRequest struct {
63+
Limit *int `json:"-" validate:"omitempty,min=0,max=1000"` // from ExtractRaw
64+
}
65+
66+
func (req *periodicJobListRequest) ExtractRaw(r *http.Request) error {
67+
if limitStr := r.URL.Query().Get("limit"); limitStr != "" {
68+
limit, err := strconv.Atoi(limitStr)
69+
if err != nil {
70+
return apierror.NewBadRequestf("Couldn't convert `limit` to integer: %s.", err)
71+
}
72+
73+
req.Limit = &limit
74+
}
75+
76+
return nil
77+
}
78+
79+
func (a *periodicJobListEndpoint[TTx]) Execute(ctx context.Context, req *periodicJobListRequest) (*listResponse[uitype.RiverPeriodicJob], error) {
80+
result, err := a.DB.PeriodicJobGetAll(ctx, &riverprodriver.PeriodicJobGetAllParams{
81+
Max: ptrutil.ValOrDefault(req.Limit, 100),
82+
Schema: a.Client.Schema(),
83+
})
84+
if err != nil {
85+
return nil, fmt.Errorf("error listing periodic jobs: %w", err)
86+
}
87+
88+
return listResponseFrom(sliceutil.Map(result, internalPeriodicJobToSerializablePeriodicJob)), nil
89+
}
90+
91+
func internalPeriodicJobToSerializablePeriodicJob(internal *riverprodriver.PeriodicJob) *uitype.RiverPeriodicJob {
92+
return (*uitype.RiverPeriodicJob)(internal)
93+
}
94+
4295
//
4396
// producerListEndpoint
4497
//

riverproui/internal/prohandler/pro_handler_api_endpoints_test.go

Lines changed: 37 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -28,11 +28,12 @@ import (
2828
"riverqueue.com/riverui/internal/riverinternaltest"
2929
"riverqueue.com/riverui/internal/riverinternaltest/testfactory"
3030
"riverqueue.com/riverui/internal/uicommontest"
31+
"riverqueue.com/riverui/riverproui/internal/protestfactory"
3132
)
3233

3334
type setupEndpointTestBundle struct {
3435
client *riverpro.Client[pgx.Tx]
35-
exec riverdriver.ExecutorTx
36+
exec driver.ProExecutorTx
3637
logger *slog.Logger
3738
tx pgx.Tx
3839
}
@@ -101,6 +102,41 @@ func testMountOpts(t *testing.T) *apiendpoint.MountOpts {
101102
}
102103
}
103104

105+
func TestProAPIHandlerPeriodicJobList(t *testing.T) {
106+
t.Parallel()
107+
108+
ctx := context.Background()
109+
110+
t.Run("Success", func(t *testing.T) {
111+
t.Parallel()
112+
113+
endpoint, bundle := setupEndpoint(ctx, t, NewPeriodicJobListEndpoint)
114+
115+
job1 := protestfactory.PeriodicJob(ctx, t, bundle.exec, &protestfactory.PeriodicJobOpts{ID: ptrutil.Ptr("alpha"), NextRunAt: ptrutil.Ptr(time.Now().Add(time.Minute))})
116+
job2 := protestfactory.PeriodicJob(ctx, t, bundle.exec, &protestfactory.PeriodicJobOpts{ID: ptrutil.Ptr("beta"), NextRunAt: ptrutil.Ptr(time.Now().Add(2 * time.Minute))})
117+
118+
resp, err := apitest.InvokeHandler(ctx, endpoint.Execute, testMountOpts(t), &periodicJobListRequest{})
119+
require.NoError(t, err)
120+
require.Len(t, resp.Data, 2)
121+
require.Equal(t, job1.ID, resp.Data[0].ID)
122+
require.Equal(t, job2.ID, resp.Data[1].ID)
123+
})
124+
125+
t.Run("Limit", func(t *testing.T) {
126+
t.Parallel()
127+
128+
endpoint, bundle := setupEndpoint(ctx, t, NewPeriodicJobListEndpoint)
129+
130+
job1 := protestfactory.PeriodicJob(ctx, t, bundle.exec, nil)
131+
_ = protestfactory.PeriodicJob(ctx, t, bundle.exec, nil)
132+
133+
resp, err := apitest.InvokeHandler(ctx, endpoint.Execute, testMountOpts(t), &periodicJobListRequest{Limit: ptrutil.Ptr(1)})
134+
require.NoError(t, err)
135+
require.Len(t, resp.Data, 1)
136+
require.Equal(t, job1.ID, resp.Data[0].ID)
137+
})
138+
}
139+
104140
func TestProAPIHandlerWorkflowCancel(t *testing.T) {
105141
t.Parallel()
106142

Lines changed: 44 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,44 @@
1+
package protestfactory
2+
3+
import (
4+
"context"
5+
"fmt"
6+
"sync/atomic"
7+
"testing"
8+
"time"
9+
10+
"github.com/stretchr/testify/require"
11+
12+
"github.com/riverqueue/river/rivershared/util/ptrutil"
13+
14+
"riverqueue.com/riverpro/driver"
15+
)
16+
17+
type PeriodicJobOpts struct {
18+
ID *string
19+
NextRunAt *time.Time
20+
UpdatedAt *time.Time
21+
}
22+
23+
func PeriodicJob(ctx context.Context, tb testing.TB, exec driver.ProExecutor, opts *PeriodicJobOpts) *driver.PeriodicJob {
24+
tb.Helper()
25+
26+
if opts == nil {
27+
opts = &PeriodicJobOpts{}
28+
}
29+
30+
periodicJob, err := exec.PeriodicJobInsert(ctx, &driver.PeriodicJobInsertParams{
31+
ID: ptrutil.ValOrDefaultFunc(opts.ID, func() string { return fmt.Sprintf("periodic_job_%05d", nextSeq()) }),
32+
NextRunAt: ptrutil.ValOrDefaultFunc(opts.NextRunAt, time.Now),
33+
UpdatedAt: opts.UpdatedAt,
34+
Schema: "",
35+
})
36+
require.NoError(tb, err)
37+
return periodicJob
38+
}
39+
40+
var seq int64 = 1 //nolint:gochecknoglobals
41+
42+
func nextSeq() int {
43+
return int(atomic.AddInt64(&seq, 1))
44+
}

riverproui/internal/uitype/ui_api_types.go

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,13 @@ type PartitionConfig struct {
1313
ByKind bool `json:"by_kind"`
1414
}
1515

16+
type RiverPeriodicJob struct {
17+
ID string `json:"id"`
18+
CreatedAt time.Time `json:"created_at"`
19+
NextRunAt time.Time `json:"next_run_at"`
20+
UpdatedAt time.Time `json:"updated_at"`
21+
}
22+
1623
type RiverProducer struct {
1724
ID int64 `json:"id"`
1825
ClientID string `json:"client_id"`

riverproui/pro_handler_test.go

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,13 +17,15 @@ import (
1717
"github.com/riverqueue/river/riverdriver"
1818

1919
"riverqueue.com/riverpro"
20+
"riverqueue.com/riverpro/driver"
2021
"riverqueue.com/riverpro/driver/riverpropgxv5"
2122

2223
"riverqueue.com/riverui"
2324
"riverqueue.com/riverui/internal/handlertest"
2425
"riverqueue.com/riverui/internal/riverinternaltest"
2526
"riverqueue.com/riverui/internal/riverinternaltest/testfactory"
2627
"riverqueue.com/riverui/internal/uicommontest"
28+
"riverqueue.com/riverui/riverproui/internal/protestfactory"
2729
"riverqueue.com/riverui/uiendpoints"
2830
)
2931

@@ -71,6 +73,11 @@ func TestProHandlerIntegration(t *testing.T) {
7173
testRunner := func(exec riverdriver.Executor, makeAPICall handlertest.APICallFunc) {
7274
ctx := context.Background()
7375

76+
proExec, ok := exec.(driver.ProExecutor)
77+
require.True(t, ok)
78+
79+
_ = protestfactory.PeriodicJob(ctx, t, proExec, nil)
80+
7481
queue := testfactory.Queue(ctx, t, exec, nil)
7582

7683
workflowID := uuid.New()
@@ -81,6 +88,7 @@ func TestProHandlerIntegration(t *testing.T) {
8188
// Verify OSS features endpoint is mounted and returns success even w/ Pro bundle:
8289
makeAPICall(t, "FeaturesGet", http.MethodGet, "/api/features", nil)
8390

91+
makeAPICall(t, "PeriodicJobList", http.MethodGet, "/api/pro/periodic-jobs", nil)
8492
makeAPICall(t, "ProducerList", http.MethodGet, "/api/pro/producers?queue_name="+queue.Name, nil)
8593
makeAPICall(t, "WorkflowCancel", http.MethodPost, fmt.Sprintf("/api/pro/workflows/%s/cancel", workflowID), nil)
8694
makeAPICall(t, "WorkflowGet", http.MethodGet, fmt.Sprintf("/api/pro/workflows/%s", workflowID2), nil)

0 commit comments

Comments
 (0)