mirror of
https://github.com/coder/coder.git
synced 2026-06-04 13:38:21 +00:00
bddb808b25
Fixes all our Go file imports to match the preferred spec that we've _mostly_ been using. For example: ``` import ( "context" "time" "github.com/prometheus/client_golang/prometheus" "golang.org/x/xerrors" "gopkg.in/natefinch/lumberjack.v2" "cdr.dev/slog/v3" "github.com/coder/coder/v2/codersdk/agentsdk" "github.com/coder/serpent" ) ``` 3 groups: standard library, 3rd partly libs, Coder libs. This PR makes the change across the codebase. The PR in the stack above modifies our formatting to maintain this state of affairs, and is a separate PR so it's possible to review that one in detail.
74 lines
1.4 KiB
Go
74 lines
1.4 KiB
Go
package tailnet
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"sync"
|
|
|
|
"github.com/google/uuid"
|
|
|
|
"cdr.dev/slog/v3"
|
|
"github.com/coder/coder/v2/coderd/database/pubsub"
|
|
)
|
|
|
|
type readyForHandshake struct {
|
|
src uuid.UUID
|
|
dst uuid.UUID
|
|
}
|
|
|
|
type handshaker struct {
|
|
ctx context.Context
|
|
logger slog.Logger
|
|
coordinatorID uuid.UUID
|
|
pubsub pubsub.Pubsub
|
|
updates <-chan readyForHandshake
|
|
|
|
workerWG sync.WaitGroup
|
|
}
|
|
|
|
func newHandshaker(ctx context.Context,
|
|
logger slog.Logger,
|
|
id uuid.UUID,
|
|
ps pubsub.Pubsub,
|
|
updates <-chan readyForHandshake,
|
|
startWorkers <-chan struct{},
|
|
) *handshaker {
|
|
s := &handshaker{
|
|
ctx: ctx,
|
|
logger: logger,
|
|
coordinatorID: id,
|
|
pubsub: ps,
|
|
updates: updates,
|
|
}
|
|
// add to the waitgroup immediately to avoid any races waiting for it before
|
|
// the workers start.
|
|
s.workerWG.Add(numHandshakerWorkers)
|
|
go func() {
|
|
<-startWorkers
|
|
for i := 0; i < numHandshakerWorkers; i++ {
|
|
go s.worker()
|
|
}
|
|
}()
|
|
return s
|
|
}
|
|
|
|
func (t *handshaker) worker() {
|
|
defer t.workerWG.Done()
|
|
|
|
for {
|
|
select {
|
|
case <-t.ctx.Done():
|
|
t.logger.Debug(t.ctx, "handshaker worker exiting", slog.Error(t.ctx.Err()))
|
|
return
|
|
|
|
case rfh := <-t.updates:
|
|
err := t.pubsub.Publish(eventReadyForHandshake, []byte(fmt.Sprintf(
|
|
"%s,%s", rfh.dst.String(), rfh.src.String(),
|
|
)))
|
|
if err != nil {
|
|
t.logger.Error(t.ctx, "publish ready for handshake", slog.Error(err))
|
|
}
|
|
}
|
|
}
|
|
}
|