Skip to content

Commit d29ce18

Browse files
committed
refactor: fix all gofmt issues and further reduce daemon cyclomatic complexity
1 parent 43d7766 commit d29ce18

4 files changed

Lines changed: 94 additions & 77 deletions

File tree

cmd/conch/main.go

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -260,7 +260,6 @@ func runElectReadOnly(ctx context.Context, logger *slog.Logger, endpoints []stri
260260
}
261261
}
262262

263-
264263
func handleSema(args []string) {
265264
fs := flag.NewFlagSet("sema", flag.ExitOnError)
266265

internal/cron/cron.go

Lines changed: 45 additions & 38 deletions
Original file line numberDiff line numberDiff line change
@@ -595,49 +595,56 @@ func (cd *Conchd) Run(ctx context.Context) error {
595595
return fmt.Errorf("watch channel closed")
596596
}
597597
for _, ev := range wresp.Events {
598-
name := strings.TrimPrefix(string(ev.Kv.Key), core.CronJobPrefix())
599-
if ev.Type == clientv3.EventTypeDelete {
600-
if cancel, ok := cancels[name]; ok {
601-
cancel()
602-
delete(cancels, name)
603-
}
604-
delete(jobs, name)
605-
cd.logger.Info("job removed", "name", name)
606-
} else {
607-
var spec JobSpec
608-
if err := json.Unmarshal(ev.Kv.Value, &spec); err == nil {
609-
oldSpec, exists := jobs[name]
610-
if exists && oldSpec.Schedule == spec.Schedule && oldSpec.RunTTL == spec.RunTTL && len(oldSpec.Cmd) == len(spec.Cmd) {
611-
cmdChanged := false
612-
for i := range spec.Cmd {
613-
if spec.Cmd[i] != oldSpec.Cmd[i] {
614-
cmdChanged = true
615-
break
616-
}
617-
}
618-
if !cmdChanged {
619-
// Execution spec is unchanged. Update jobs map for metadata, but do NOT restart scheduler.
620-
jobs[name] = spec
621-
cd.logger.Debug("job touched but execution spec unchanged; not restarting scheduler", "name", name)
622-
continue
623-
}
624-
}
625-
626-
if cancel, ok := cancels[name]; ok {
627-
cancel()
628-
}
629-
jobs[name] = spec
630-
jobCtx, cancel := context.WithCancel(ctx)
631-
cancels[name] = cancel
632-
go runJobScheduler(ctx, jobCtx, cd.logger, cd.client, sess, name, spec)
633-
cd.logger.Info("job added/updated", "name", name)
634-
}
635-
}
598+
cd.handleWatchEvent(ctx, ev, jobs, cancels, sess)
636599
}
637600
}
638601
}
639602
}
640603

604+
func (cd *Conchd) handleWatchEvent(ctx context.Context, ev *clientv3.Event, jobs map[string]JobSpec, cancels map[string]context.CancelFunc, sess *core.CoreSession) {
605+
name := strings.TrimPrefix(string(ev.Kv.Key), core.CronJobPrefix())
606+
if ev.Type == clientv3.EventTypeDelete {
607+
if cancel, ok := cancels[name]; ok {
608+
cancel()
609+
delete(cancels, name)
610+
}
611+
delete(jobs, name)
612+
cd.logger.Info("job removed", "name", name)
613+
return
614+
}
615+
616+
var spec JobSpec
617+
if err := json.Unmarshal(ev.Kv.Value, &spec); err != nil {
618+
return
619+
}
620+
621+
oldSpec, exists := jobs[name]
622+
if exists && oldSpec.Schedule == spec.Schedule && oldSpec.RunTTL == spec.RunTTL && len(oldSpec.Cmd) == len(spec.Cmd) {
623+
cmdChanged := false
624+
for i := range spec.Cmd {
625+
if spec.Cmd[i] != oldSpec.Cmd[i] {
626+
cmdChanged = true
627+
break
628+
}
629+
}
630+
if !cmdChanged {
631+
// Execution spec is unchanged. Update jobs map for metadata, but do NOT restart scheduler.
632+
jobs[name] = spec
633+
cd.logger.Debug("job touched but execution spec unchanged; not restarting scheduler", "name", name)
634+
return
635+
}
636+
}
637+
638+
if cancel, ok := cancels[name]; ok {
639+
cancel()
640+
}
641+
jobs[name] = spec
642+
jobCtx, cancel := context.WithCancel(ctx)
643+
cancels[name] = cancel
644+
go runJobScheduler(ctx, jobCtx, cd.logger, cd.client, sess, name, spec)
645+
cd.logger.Info("job added/updated", "name", name)
646+
}
647+
641648
func runJobScheduler(daemonCtx, jobCtx context.Context, logger *slog.Logger, client *clientv3.Client, sess *core.CoreSession, name string, spec JobSpec) {
642649
parser := cron.NewParser(cron.Minute | cron.Hour | cron.Dom | cron.Month | cron.Dow | cron.Descriptor)
643650
sched, err := parser.Parse(spec.Schedule)

internal/elect/elect.go

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -346,4 +346,3 @@ func updateCandidates(candidates []candidate, events []*clientv3.Event) []candid
346346
}
347347
return candidates
348348
}
349-

internal/sema/sema.go

Lines changed: 49 additions & 37 deletions
Original file line numberDiff line numberDiff line change
@@ -106,56 +106,68 @@ func (sh *SemaHolder) acquireInternal(ctx context.Context) (int64, error) {
106106
prefix := core.SemaPrefix(sh.name, sh.max)
107107

108108
for {
109-
// 2. Fetch all keys under the prefix sorted by create-revision
110-
getResp, err := sh.client.Get(ctx, prefix,
111-
clientv3.WithPrefix(),
112-
clientv3.WithSort(clientv3.SortByCreateRevision, clientv3.SortAscend),
113-
)
109+
rev, acquired, err := sh.checkRankAndWatch(ctx, prefix)
114110
if err != nil {
115111
return 0, err
116112
}
117-
118-
// Find our rank
119-
rank := -1
120-
for i, kv := range getResp.Kvs {
121-
if string(kv.Key) == sh.key {
122-
rank = i
123-
break
124-
}
113+
if acquired {
114+
return rev, nil
125115
}
116+
}
117+
}
118+
}
126119

127-
if rank == -1 {
128-
return 0, fmt.Errorf("our key was deleted from etcd")
129-
}
120+
func (sh *SemaHolder) checkRankAndWatch(ctx context.Context, prefix string) (int64, bool, error) {
121+
// 2. Fetch all keys under the prefix sorted by create-revision
122+
getResp, err := sh.client.Get(ctx, prefix,
123+
clientv3.WithPrefix(),
124+
clientv3.WithSort(clientv3.SortByCreateRevision, clientv3.SortAscend),
125+
)
126+
if err != nil {
127+
return 0, false, err
128+
}
130129

131-
// 3. If rank < N, we hold a slot!
132-
if rank < sh.max {
133-
sh.acquired = true
134-
return getResp.Kvs[rank].CreateRevision, nil
135-
}
130+
// Find our rank
131+
rank := -1
132+
for i, kv := range getResp.Kvs {
133+
if string(kv.Key) == sh.key {
134+
rank = i
135+
break
136+
}
137+
}
136138

137-
// Non-blocking means we don't wait if rank >= max
138-
if sh.waitLimit == 0 {
139-
return 0, context.DeadlineExceeded // Treat non-block as immediate timeout
140-
}
139+
if rank == -1 {
140+
return 0, false, fmt.Errorf("our key was deleted from etcd")
141+
}
141142

142-
// 4. Watch for deletion of the key at rank - max
143-
targetKey := string(getResp.Kvs[rank-sh.max].Key)
143+
// 3. If rank < N, we hold a slot!
144+
if rank < sh.max {
145+
sh.acquired = true
146+
return getResp.Kvs[rank].CreateRevision, true, nil
147+
}
144148

145-
// Watch from the revision of our Get response to avoid races
146-
err = sh.watchKeyDeletion(ctx, targetKey, getResp.Header.Revision, "session lost while waiting for semaphore slot")
147-
if err != nil {
148-
return 0, err
149-
}
149+
// Non-blocking means we don't wait if rank >= max
150+
if sh.waitLimit == 0 {
151+
return 0, false, context.DeadlineExceeded // Treat non-block as immediate timeout
152+
}
150153

151-
// Check context done
152-
if ctx.Err() != nil {
153-
return 0, ctx.Err()
154-
}
155-
}
154+
// 4. Watch for deletion of the key at rank - max
155+
targetKey := string(getResp.Kvs[rank-sh.max].Key)
156+
157+
// Watch from the revision of our Get response to avoid races
158+
err = sh.watchKeyDeletion(ctx, targetKey, getResp.Header.Revision, "session lost while waiting for semaphore slot")
159+
if err != nil {
160+
return 0, false, err
156161
}
162+
163+
// Check context done
164+
if ctx.Err() != nil {
165+
return 0, false, ctx.Err()
166+
}
167+
return 0, false, nil
157168
}
158169

170+
159171
func (sh *SemaHolder) watchKeyDeletion(ctx context.Context, targetKey string, startRev int64, sessionErrMsg string) error {
160172
watchCtx, watchCancel := context.WithCancel(ctx)
161173
defer watchCancel()

0 commit comments

Comments
 (0)