mirror of
https://github.com/coder/coder.git
synced 2026-09-24 15:04:27 +08:00
feat: Refactor API routes to use UUIDs instead of friendly names (#401)
* Add client for agent * Cleanup code * Fix linting error * Rename routes to be simpler * Rename workspace history to workspace build * Refactor HTTP middlewares to use UUIDs * Cleanup routes * Compiles! * Fix files and organizations * Fix querying * Fix agent lock * Cleanup database abstraction * Add parameters * Fix linting errors * Fix log race * Lock on close wait * Fix log cleanup * Fix e2e tests * Fix upstream version of opencensus-go * Update coderdtest.go * Fix coverpkg * Fix codecov ignore
This commit is contained in:
@@ -30,8 +30,9 @@ func Listen(connListener net.Listener, iceServersFunc ICEServersFunc, opts *peer
|
||||
}
|
||||
ctx, cancelFunc := context.WithCancel(context.Background())
|
||||
listener := &Listener{
|
||||
connectionChannel: make(chan *peer.Conn),
|
||||
iceServersFunc: iceServersFunc,
|
||||
connectionChannel: make(chan *peer.Conn),
|
||||
connectionListener: connListener,
|
||||
iceServersFunc: iceServersFunc,
|
||||
|
||||
closeFunc: cancelFunc,
|
||||
closed: make(chan struct{}),
|
||||
@@ -56,8 +57,9 @@ func Listen(connListener net.Listener, iceServersFunc ICEServersFunc, opts *peer
|
||||
}
|
||||
|
||||
type Listener struct {
|
||||
connectionChannel chan *peer.Conn
|
||||
iceServersFunc ICEServersFunc
|
||||
connectionChannel chan *peer.Conn
|
||||
connectionListener net.Listener
|
||||
iceServersFunc ICEServersFunc
|
||||
|
||||
closeFunc context.CancelFunc
|
||||
closed chan struct{}
|
||||
@@ -89,6 +91,7 @@ func (l *Listener) closeWithError(err error) error {
|
||||
return l.closeError
|
||||
}
|
||||
|
||||
_ = l.connectionListener.Close()
|
||||
l.closeError = err
|
||||
l.closeFunc()
|
||||
close(l.closed)
|
||||
|
||||
+25
-5
@@ -2,6 +2,7 @@ package peerbroker
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/base64"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
@@ -118,7 +119,7 @@ func (p *proxyListen) NegotiateConnection(stream proto.DRPCPeerBroker_NegotiateC
|
||||
return xerrors.Errorf("maximum payload size %d exceeded", maxPayloadSizeBytes)
|
||||
}
|
||||
data = append([]byte(streamID), data...)
|
||||
err = p.pubsub.Publish(proxyOutID(p.channelID), data)
|
||||
err = p.pubsub.Publish(proxyOutID(p.channelID), marshal(data))
|
||||
if err != nil {
|
||||
return xerrors.Errorf("publish: %w", err)
|
||||
}
|
||||
@@ -127,6 +128,11 @@ func (p *proxyListen) NegotiateConnection(stream proto.DRPCPeerBroker_NegotiateC
|
||||
}
|
||||
|
||||
func (*proxyListen) onServerToClientMessage(streamID string, stream proto.DRPCPeerBroker_NegotiateConnectionStream, message []byte) error {
|
||||
var err error
|
||||
message, err = unmarshal(message)
|
||||
if err != nil {
|
||||
return xerrors.Errorf("decode: %w", err)
|
||||
}
|
||||
if len(message) < streamIDLength {
|
||||
return xerrors.Errorf("got message length %d < %d", len(message), streamIDLength)
|
||||
}
|
||||
@@ -136,7 +142,7 @@ func (*proxyListen) onServerToClientMessage(streamID string, stream proto.DRPCPe
|
||||
return nil
|
||||
}
|
||||
var msg proto.Exchange
|
||||
err := protobuf.Unmarshal(message[streamIDLength:], &msg)
|
||||
err = protobuf.Unmarshal(message[streamIDLength:], &msg)
|
||||
if err != nil {
|
||||
return xerrors.Errorf("unmarshal message: %w", err)
|
||||
}
|
||||
@@ -173,10 +179,14 @@ func (p *proxyDial) listen() error {
|
||||
}
|
||||
|
||||
func (p *proxyDial) onClientToServerMessage(ctx context.Context, message []byte) error {
|
||||
var err error
|
||||
message, err = unmarshal(message)
|
||||
if err != nil {
|
||||
return xerrors.Errorf("decode: %w", err)
|
||||
}
|
||||
if len(message) < streamIDLength {
|
||||
return xerrors.Errorf("got message length %d < %d", len(message), streamIDLength)
|
||||
}
|
||||
var err error
|
||||
streamID := string(message[0:streamIDLength])
|
||||
p.streamMutex.Lock()
|
||||
stream, ok := p.streams[streamID]
|
||||
@@ -190,7 +200,7 @@ func (p *proxyDial) onClientToServerMessage(ctx context.Context, message []byte)
|
||||
go func() {
|
||||
defer stream.Close()
|
||||
|
||||
err = p.onServerToClientMessage(streamID, stream)
|
||||
err := p.onServerToClientMessage(streamID, stream)
|
||||
if err != nil {
|
||||
p.logger.Debug(ctx, "failed to accept server message", slog.Error(err))
|
||||
}
|
||||
@@ -236,7 +246,7 @@ func (p *proxyDial) onServerToClientMessage(streamID string, stream proto.DRPCPe
|
||||
return xerrors.Errorf("maximum payload size %d exceeded", maxPayloadSizeBytes)
|
||||
}
|
||||
data = append([]byte(streamID), data...)
|
||||
err = p.pubsub.Publish(proxyInID(p.channelID), data)
|
||||
err = p.pubsub.Publish(proxyInID(p.channelID), marshal(data))
|
||||
if err != nil {
|
||||
return xerrors.Errorf("publish: %w", err)
|
||||
}
|
||||
@@ -251,6 +261,16 @@ func (p *proxyDial) Close() error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// base64 needs to be used here to keep the pubsub messages in UTF-8 range.
|
||||
// PostgreSQL cannot handle non UTF-8 messages over pubsub.
|
||||
func marshal(data []byte) []byte {
|
||||
return []byte(base64.StdEncoding.EncodeToString(data))
|
||||
}
|
||||
|
||||
func unmarshal(data []byte) ([]byte, error) {
|
||||
return base64.StdEncoding.DecodeString(string(data))
|
||||
}
|
||||
|
||||
func proxyOutID(channelID string) string {
|
||||
return fmt.Sprintf("%s-out", channelID)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user