Skip to content

Commit 48ca2fc

Browse files
committed
refactor: use generics to eliminate rudundant logic
1 parent f2f34fa commit 48ca2fc

113 files changed

Lines changed: 2100 additions & 5098 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

.github/.golangci.yml

Lines changed: 4 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
version: "2"
1+
version: '2'
22
run:
33
tests: false
44
timeout: 5m
@@ -11,7 +11,6 @@ linters:
1111
- asciicheck
1212
- bidichk
1313
- bodyclose
14-
- canonicalheader
1514
- clickhouselint
1615
- contextcheck
1716
- copyloopvar
@@ -114,13 +113,13 @@ linters:
114113
rules:
115114
- linters:
116115
- revive
117-
text: "var-naming:"
116+
text: 'var-naming:'
118117
- linters:
119118
- revive
120-
text: "unused-parameter:"
119+
text: 'unused-parameter:'
121120
- linters:
122121
- revive
123-
text: "exported:"
122+
text: 'exported:'
124123
- linters:
125124
- forbidigo
126125
path: cli/

.vscode/settings.json

Lines changed: 12 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,5 +22,16 @@
2222
},
2323
"tasks.statusbar.default.hide": true,
2424
"tasks.statusbar.limit": 8,
25-
"js/ts.tsdk.path": "node_modules/@typescript/native-preview"
25+
"js/ts.tsdk.path": "node_modules/@typescript/native-preview",
26+
"yaml.disableSchemaDetection": [
27+
"**/.github/workflows/*.yml",
28+
"**/.github/workflows/*.yaml",
29+
"**/.gitea/workflows/*.yml",
30+
"**/.gitea/workflows/*.yaml",
31+
"**/.forgejo/workflows/*.yml",
32+
"**/.forgejo/workflows/*.yaml"
33+
],
34+
"yaml.schemas": {
35+
"https://json.schemastore.org/github-workflow.json": ".github/workflows/*"
36+
}
2637
}

backend/internal/activity/handler.go

Lines changed: 22 additions & 38 deletions
Original file line numberDiff line numberDiff line change
@@ -56,28 +56,16 @@ type ListActivitiesInput struct {
5656
ResourceType string `query:"resourceType" doc:"Filter by resource type"`
5757
}
5858

59-
type ListActivitiesOutput struct {
60-
Body base.Paginated[activitytypes.Activity]
61-
}
62-
6359
type GetActivityInput struct {
6460
EnvironmentID string `path:"id" doc:"Environment ID"`
6561
ActivityID string `path:"activityId" doc:"Activity ID"`
6662
Limit int `query:"limit" default:"500" doc:"Maximum messages to return"`
6763
}
6864

69-
type GetActivityOutput struct {
70-
Body base.ApiResponse[activitytypes.Detail]
71-
}
72-
7365
type ClearActivityHistoryInput struct {
7466
EnvironmentID string `path:"id" doc:"Environment ID"`
7567
}
7668

77-
type ClearActivityHistoryOutput struct {
78-
Body base.ApiResponse[activitytypes.ClearHistoryResult]
79-
}
80-
8169
type StreamAllActivitiesInput struct {
8270
Limit int `query:"limit" default:"50" doc:"Snapshot limit per environment"`
8371
}
@@ -95,10 +83,6 @@ type CancelActivityInput struct {
9583
RequestedBy string `query:"requestedBy" doc:"Display name to attribute the cancellation to (used when proxying to a remote environment)"`
9684
}
9785

98-
type CancelActivityOutput struct {
99-
Body base.ApiResponse[activitytypes.Activity]
100-
}
101-
10286
func NewHandler(activityService *ActivityService, environment EnvironmentDependencies) *ActivityHandler {
10387
return &ActivityHandler{
10488
activityService: activityService,
@@ -149,7 +133,7 @@ func RegisterActivities(api huma.API, h *ActivityHandler) {
149133
}, authz.PermActivitiesDelete, h.ClearHistory)
150134
}
151135

152-
func (h *ActivityHandler) ListActivities(ctx context.Context, input *ListActivitiesInput) (*ListActivitiesOutput, error) {
136+
func (h *ActivityHandler) ListActivities(ctx context.Context, input *ListActivitiesInput) (*handlerutil.Page[activitytypes.Activity], error) {
153137
if input.EnvironmentID != "0" {
154138
return h.proxyListActivitiesInternal(ctx, input)
155139
}
@@ -171,7 +155,7 @@ func (h *ActivityHandler) ListActivities(ctx context.Context, input *ListActivit
171155
}
172156
h.applyActivitySourceLabelsInternal(ctx, input.EnvironmentID, activities)
173157

174-
return &ListActivitiesOutput{
158+
return &handlerutil.Page[activitytypes.Activity]{
175159
Body: base.Paginated[activitytypes.Activity]{
176160
Success: true,
177161
Data: activities,
@@ -180,7 +164,7 @@ func (h *ActivityHandler) ListActivities(ctx context.Context, input *ListActivit
180164
}, nil
181165
}
182166

183-
func (h *ActivityHandler) GetActivity(ctx context.Context, input *GetActivityInput) (*GetActivityOutput, error) {
167+
func (h *ActivityHandler) GetActivity(ctx context.Context, input *GetActivityInput) (*handlerutil.Out[activitytypes.Detail], error) {
184168
if input.EnvironmentID != "0" {
185169
return h.proxyGetActivityInternal(ctx, input)
186170
}
@@ -197,15 +181,15 @@ func (h *ActivityHandler) GetActivity(ctx context.Context, input *GetActivityInp
197181
}
198182
h.applyActivitySourceLabelInternal(ctx, input.EnvironmentID, &detail.Activity)
199183

200-
return &GetActivityOutput{
184+
return &handlerutil.Out[activitytypes.Detail]{
201185
Body: base.ApiResponse[activitytypes.Detail]{
202186
Success: true,
203187
Data: *detail,
204188
},
205189
}, nil
206190
}
207191

208-
func (h *ActivityHandler) ClearHistory(ctx context.Context, input *ClearActivityHistoryInput) (*ClearActivityHistoryOutput, error) {
192+
func (h *ActivityHandler) ClearHistory(ctx context.Context, input *ClearActivityHistoryInput) (*handlerutil.Out[activitytypes.ClearHistoryResult], error) {
209193
if input.EnvironmentID != "0" {
210194
return h.proxyClearHistoryInternal(ctx, input)
211195
}
@@ -215,15 +199,15 @@ func (h *ActivityHandler) ClearHistory(ctx context.Context, input *ClearActivity
215199
return nil, huma.Error500InternalServerError(err.Error())
216200
}
217201

218-
return &ClearActivityHistoryOutput{
202+
return &handlerutil.Out[activitytypes.ClearHistoryResult]{
219203
Body: base.ApiResponse[activitytypes.ClearHistoryResult]{
220204
Success: true,
221205
Data: activitytypes.ClearHistoryResult{Deleted: deleted},
222206
},
223207
}, nil
224208
}
225209

226-
func (h *ActivityHandler) CancelActivity(ctx context.Context, input *CancelActivityInput) (*CancelActivityOutput, error) {
210+
func (h *ActivityHandler) CancelActivity(ctx context.Context, input *CancelActivityInput) (*handlerutil.Out[activitytypes.Activity], error) {
227211
if input.EnvironmentID != "0" {
228212
return h.proxyCancelActivityInternal(ctx, input)
229213
}
@@ -245,25 +229,25 @@ func (h *ActivityHandler) CancelActivity(ctx context.Context, input *CancelActiv
245229
}
246230
h.applyActivitySourceLabelInternal(ctx, input.EnvironmentID, cancelled)
247231

248-
return &CancelActivityOutput{
232+
return &handlerutil.Out[activitytypes.Activity]{
249233
Body: base.ApiResponse[activitytypes.Activity]{
250234
Success: true,
251235
Data: *cancelled,
252236
},
253237
}, nil
254238
}
255239

256-
func (h *ActivityHandler) proxyCancelActivityInternal(ctx context.Context, input *CancelActivityInput) (*CancelActivityOutput, error) {
240+
func (h *ActivityHandler) proxyCancelActivityInternal(ctx context.Context, input *CancelActivityInput) (*handlerutil.Out[activitytypes.Activity], error) {
257241
path := fmt.Sprintf("/api/environments/0/activities/%s/cancel", url.PathEscape(input.ActivityID))
258242
if requestedBy := h.cancelRequestedByInternal(ctx, input.RequestedBy); requestedBy != "" {
259243
path += "?requestedBy=" + url.QueryEscape(requestedBy)
260244
}
261-
out, err := handlerutil.ProxyRemoteJSON[base.ApiResponse[activitytypes.Activity]](ctx, h.environment.ProxyJSONRequest, input.EnvironmentID, http.MethodPost, path, nil)
245+
out, err := h.environment.ProxyJSONRequest.JSON[base.ApiResponse[activitytypes.Activity]](ctx, input.EnvironmentID, http.MethodPost, path, nil)
262246
if err != nil {
263247
return nil, err
264248
}
265249
h.applyActivitySourceLabelInternal(ctx, input.EnvironmentID, &out.Data)
266-
return &CancelActivityOutput{Body: *out}, nil
250+
return &handlerutil.Out[activitytypes.Activity]{Body: *out}, nil
267251
}
268252

269253
// cancelRequestedByInternal resolves a human-readable name for the cancellation
@@ -480,42 +464,42 @@ func (h *ActivityHandler) runRemoteActivityStreamPollerInternal(ctx context.Cont
480464
}
481465
}
482466

483-
func (h *ActivityHandler) proxyListActivitiesInternal(ctx context.Context, input *ListActivitiesInput) (*ListActivitiesOutput, error) {
467+
func (h *ActivityHandler) proxyListActivitiesInternal(ctx context.Context, input *ListActivitiesInput) (*handlerutil.Page[activitytypes.Activity], error) {
484468
path := "/api/environments/0/activities?" + activityListQueryInternal(input).Encode()
485-
out, err := handlerutil.ProxyRemoteJSON[base.Paginated[activitytypes.Activity]](ctx, h.environment.ProxyJSONRequest, input.EnvironmentID, http.MethodGet, path, nil)
469+
out, err := h.environment.ProxyJSONRequest.JSON[base.Paginated[activitytypes.Activity]](ctx, input.EnvironmentID, http.MethodGet, path, nil)
486470
if err != nil {
487471
return nil, err
488472
}
489473
h.applyActivitySourceLabelsInternal(ctx, input.EnvironmentID, out.Data)
490-
return &ListActivitiesOutput{Body: *out}, nil
474+
return &handlerutil.Page[activitytypes.Activity]{Body: *out}, nil
491475
}
492476

493-
func (h *ActivityHandler) proxyListActivitiesForEnvironmentInternal(ctx context.Context, environment environment.Environment, input *ListActivitiesInput) (*ListActivitiesOutput, error) {
477+
func (h *ActivityHandler) proxyListActivitiesForEnvironmentInternal(ctx context.Context, environment environment.Environment, input *ListActivitiesInput) (*handlerutil.Page[activitytypes.Activity], error) {
494478
path := "/api/environments/0/activities?" + activityListQueryInternal(input).Encode()
495479
var out base.Paginated[activitytypes.Activity]
496480
if err := h.environment.ProxyJSONRequestForEnvironment(ctx, environment, http.MethodGet, path, nil, &out); err != nil {
497481
return nil, handlerutil.TranslateRemoteProxyError(err)
498482
}
499483
applyActivitySourceLabelsForEnvironmentInternal(environment, out.Data)
500-
return &ListActivitiesOutput{Body: out}, nil
484+
return &handlerutil.Page[activitytypes.Activity]{Body: out}, nil
501485
}
502486

503-
func (h *ActivityHandler) proxyGetActivityInternal(ctx context.Context, input *GetActivityInput) (*GetActivityOutput, error) {
487+
func (h *ActivityHandler) proxyGetActivityInternal(ctx context.Context, input *GetActivityInput) (*handlerutil.Out[activitytypes.Detail], error) {
504488
path := fmt.Sprintf("/api/environments/0/activities/%s?limit=%d", url.PathEscape(input.ActivityID), input.Limit)
505-
out, err := handlerutil.ProxyRemoteJSON[base.ApiResponse[activitytypes.Detail]](ctx, h.environment.ProxyJSONRequest, input.EnvironmentID, http.MethodGet, path, nil)
489+
out, err := h.environment.ProxyJSONRequest.JSON[base.ApiResponse[activitytypes.Detail]](ctx, input.EnvironmentID, http.MethodGet, path, nil)
506490
if err != nil {
507491
return nil, err
508492
}
509493
h.applyActivitySourceLabelInternal(ctx, input.EnvironmentID, &out.Data.Activity)
510-
return &GetActivityOutput{Body: *out}, nil
494+
return &handlerutil.Out[activitytypes.Detail]{Body: *out}, nil
511495
}
512496

513-
func (h *ActivityHandler) proxyClearHistoryInternal(ctx context.Context, input *ClearActivityHistoryInput) (*ClearActivityHistoryOutput, error) {
514-
out, err := handlerutil.ProxyRemoteJSON[base.ApiResponse[activitytypes.ClearHistoryResult]](ctx, h.environment.ProxyJSONRequest, input.EnvironmentID, http.MethodDelete, "/api/environments/0/activities/history", nil)
497+
func (h *ActivityHandler) proxyClearHistoryInternal(ctx context.Context, input *ClearActivityHistoryInput) (*handlerutil.Out[activitytypes.ClearHistoryResult], error) {
498+
out, err := h.environment.ProxyJSONRequest.JSON[base.ApiResponse[activitytypes.ClearHistoryResult]](ctx, input.EnvironmentID, http.MethodDelete, "/api/environments/0/activities/history", nil)
515499
if err != nil {
516500
return nil, err
517501
}
518-
return &ClearActivityHistoryOutput{Body: *out}, nil
502+
return &handlerutil.Out[activitytypes.ClearHistoryResult]{Body: *out}, nil
519503
}
520504

521505
func (h *ActivityHandler) applyActivitySourceLabelsInternal(ctx context.Context, environmentID string, activities []activitytypes.Activity) {

backend/internal/actors/actor.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -113,12 +113,12 @@ func (a *Actor) Stop(ctx context.Context) error {
113113
}
114114

115115
// Request sends one typed request and consumes its response.
116-
func Request[K comparable, Q, R any](ctx context.Context, a *Actor, request Message[K, Q]) (Message[K, R], error) {
116+
func (a *Actor) Request[K comparable, Q, R any](ctx context.Context, request Message[K, Q]) (Message[K, R], error) {
117117
if a == nil {
118118
var zero Message[K, R]
119119
return zero, errors.New("actor unavailable")
120120
}
121-
return requestHandleInternal[K, Q, R](ctx, a.handle, request)
121+
return a.handle.request[K, Q, R](ctx, request)
122122
}
123123

124124
func (a *behaviorActorInternal) Receive(ctx *actor.Context) {

backend/internal/actors/executor.go

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -65,9 +65,9 @@ func NewExecutor(ctx context.Context, runtime *Runtime, kind, id string, maxRest
6565

6666
// Execute queues work, waits for its typed result, and runs after only once the
6767
// result is available to the caller. Work and after are both panic-contained.
68-
func Execute[T any](ctx context.Context, executor *Executor, label string, work func(context.Context) (T, error), after func(T, error)) (T, error) {
68+
func (e *Executor) Execute[T any](ctx context.Context, label string, work func(context.Context) (T, error), after func(T, error)) (T, error) {
6969
var zero T
70-
task, err := Submit(ctx, executor, label, work, after)
70+
task, err := e.Submit(ctx, label, work, after)
7171
if err != nil {
7272
return zero, err
7373
}
@@ -77,8 +77,8 @@ func Execute[T any](ctx context.Context, executor *Executor, label string, work
7777
// Submit queues work synchronously and returns a handle that can be awaited
7878
// separately. This is useful for terminal cleanup that must stay queued even
7979
// when the shutdown wait context expires.
80-
func Submit[T any](ctx context.Context, executor *Executor, label string, work func(context.Context) (T, error), after func(T, error)) (*Task[T], error) {
81-
if executor == nil || executor.actor == nil {
80+
func (e *Executor) Submit[T any](ctx context.Context, label string, work func(context.Context) (T, error), after func(T, error)) (*Task[T], error) {
81+
if e == nil || e.actor == nil {
8282
return nil, errors.New("actor executor unavailable")
8383
}
8484
if ctx == nil {
@@ -92,7 +92,7 @@ func Submit[T any](ctx context.Context, executor *Executor, label string, work f
9292
}
9393

9494
result := NewPromise[executorResultInternal[T]]()
95-
if err := executor.actor.Send(executorTaskValueInternal[T]{
95+
if err := e.actor.Send(executorTaskValueInternal[T]{
9696
ctx: ctx,
9797
label: label,
9898
work: work,
@@ -102,7 +102,7 @@ func Submit[T any](ctx context.Context, executor *Executor, label string, work f
102102
}); err != nil {
103103
return nil, err
104104
}
105-
return &Task[T]{result: result, actor: executor.actor}, nil
105+
return &Task[T]{result: result, actor: e.actor}, nil
106106
}
107107

108108
// Wait waits for the submitted task result without changing the task lifetime.
@@ -114,7 +114,7 @@ func (t *Task[T]) Wait(ctx context.Context) (T, error) {
114114
if ctx == nil {
115115
return zero, errors.New("task wait context unavailable")
116116
}
117-
completed, err := awaitInternal(ctx, t.actor.handle, t.result, errors.New("actor executor stopped"))
117+
completed, err := t.actor.handle.await(ctx, t.result, errors.New("actor executor stopped"))
118118
if err != nil {
119119
return zero, err
120120
}

backend/internal/actors/executor_test.go

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,7 @@ func TestExecutorSerializesHeterogeneousTasksInternal(t *testing.T) {
2020
releaseFirst := make(chan struct{})
2121
firstResult := make(chan error, 1)
2222
go func() {
23-
_, executeErr := Execute(t.Context(), executor, "first", func(context.Context) (int, error) {
23+
_, executeErr := executor.Execute(t.Context(), "first", func(context.Context) (int, error) {
2424
close(firstStarted)
2525
<-releaseFirst
2626
return 1, nil
@@ -32,7 +32,7 @@ func TestExecutorSerializesHeterogeneousTasksInternal(t *testing.T) {
3232
secondStarted := make(chan struct{})
3333
secondResult := make(chan string, 1)
3434
go func() {
35-
value, executeErr := Execute(t.Context(), executor, "second", func(context.Context) (string, error) {
35+
value, executeErr := executor.Execute(t.Context(), "second", func(context.Context) (string, error) {
3636
close(secondStarted)
3737
return "two", nil
3838
}, nil)
@@ -65,7 +65,7 @@ func TestExecutorPublishesResultBeforeCompletionCallbackInternal(t *testing.T) {
6565
releaseAfter := make(chan struct{})
6666
result := make(chan int, 1)
6767
go func() {
68-
value, _ := Execute(t.Context(), executor, "ordered", func(context.Context) (int, error) {
68+
value, _ := executor.Execute(t.Context(), "ordered", func(context.Context) (int, error) {
6969
return 42, nil
7070
}, func(int, error) {
7171
close(afterStarted)
@@ -91,12 +91,12 @@ func TestExecutorContainsTaskPanicAndContinuesInternal(t *testing.T) {
9191
executor, err := NewExecutor(t.Context(), runtime, "executor-test", "panic", 1)
9292
require.NoError(t, err)
9393

94-
_, err = Execute(t.Context(), executor, "panic task", func(context.Context) (NoPayload, error) {
94+
_, err = executor.Execute(t.Context(), "panic task", func(context.Context) (NoPayload, error) {
9595
panic("deliberate executor panic")
9696
}, nil)
9797
require.ErrorContains(t, err, "panic task panicked")
9898

99-
value, err := Execute(t.Context(), executor, "following task", func(context.Context) (int, error) {
99+
value, err := executor.Execute(t.Context(), "following task", func(context.Context) (int, error) {
100100
return 7, nil
101101
}, nil)
102102
require.NoError(t, err)

backend/internal/actors/gate.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -99,7 +99,7 @@ func (g *Gate[K]) TryAcquire(ctx context.Context, key K) (*Lease[K], bool, error
9999

100100
// Admission must resolve after it is queued; abandoning an admitted reply on
101101
// caller cancellation would leak the key without a lease to release it.
102-
admitted, err := awaitInternal(context.WithoutCancel(ctx), g.actor.handle, reply, errors.New("actor gate stopped"))
102+
admitted, err := g.actor.handle.await(context.WithoutCancel(ctx), reply, errors.New("actor gate stopped"))
103103
if err != nil {
104104
return nil, false, err
105105
}
@@ -125,7 +125,7 @@ func (l *Lease[K]) Release() {
125125
}); err != nil {
126126
return
127127
}
128-
_, _ = awaitInternal(context.Background(), l.gate.actor.handle, reply, errors.New("actor gate stopped"))
128+
_, _ = l.gate.actor.handle.await(context.Background(), reply, errors.New("actor gate stopped"))
129129
})
130130
}
131131

0 commit comments

Comments
 (0)