Clean up dead code across the codebase

Spring cleaning!
A very mechanical cleanup using several linters (unused, deadcode,
structcheck). Build and tests still pass so no behavior should be
affected.
This commit is contained in:
Andrew Lytvynov
2020-04-09 21:10:12 +00:00
committed by Andrew Lytvynov
parent 24029efcfc
commit f8661edea3
47 changed files with 16 additions and 717 deletions
-14
View File
@@ -446,20 +446,6 @@ func migrateLegacyResources(cfg InitConfig, asrv *AuthServer) error {
return nil
}
// checkRules returns true if any of the provided rules contain specified resource/verbs.
func checkRules(rules []services.Rule, kind string, verbs []string) bool {
for _, rule := range rules {
if rule.HasResource(kind) || rule.HasResource(services.Wildcard) {
for _, verb := range verbs {
if rule.HasVerb(verb) || rule.HasVerb(services.Wildcard) {
return true
}
}
}
}
return false
}
// isFirstStart returns 'true' if the auth server is starting for the 1st time
// on this server.
func isFirstStart(authServer *AuthServer, cfg InitConfig) (bool, error) {
-21
View File
@@ -600,27 +600,6 @@ func startRollingBackRotation(ca services.CertAuthority) error {
return nil
}
// completeRollingBackRotation completes rollback of the rotation and sets it to the standby state
func completeRollingBackRotation(clock clockwork.Clock, ca services.CertAuthority) error {
rotation := ca.GetRotation()
// clean up the state
rotation.Started = time.Time{}
rotation.State = services.RotationStateStandby
rotation.Phase = services.RotationPhaseStandby
rotation.Mode = ""
rotation.Schedule = services.RotationSchedule{}
keyPairs := ca.GetTLSKeyPairs()
// only keep the original certificate authority as trusted
// and remove everything else.
keyPairs = []services.TLSKeyPair{keyPairs[0]}
ca.SetTLSKeyPairs(keyPairs)
ca.SetRotation(rotation)
return nil
}
// completeRotation completes the certificate authority rotation.
func completeRotation(clock clockwork.Clock, ca services.CertAuthority) error {
rotation := ca.GetRotation()
-34
View File
@@ -17,7 +17,6 @@ limitations under the License.
package auth
import (
"io"
"net"
"sync"
"time"
@@ -27,13 +26,6 @@ import (
"golang.org/x/crypto/ssh"
)
// Implements a fake "socket" (net.Listener interface) on top of existing ssh.Channel
type fakeSocket struct {
closed chan int
connections chan net.Conn
closeOnce sync.Once
}
// FakeSSHConnection implements net.Conn interface on top of the ssh.Cnahhel
// object. This allows us to run non-SSH servers (like HTTP) on top of an
// existing SSH connection
@@ -82,29 +74,3 @@ func (conn *FakeSSHConnection) SetReadDeadline(t time.Time) error {
func (conn *FakeSSHConnection) SetWriteDeadline(t time.Time) error {
return nil
}
// Accept waits for new connections to arrive (via CreateBridge) and returns them to
// the blocked http.Serve()
func (socket *fakeSocket) Accept() (c net.Conn, err error) {
select {
case newConnection := <-socket.connections:
return newConnection, nil
case <-socket.closed:
return nil, io.EOF
}
}
// Close closes the listener.
// Any blocked Accept operations will be unblocked and return errors.
func (socket *fakeSocket) Close() error {
socket.closeOnce.Do(func() {
// broadcast that listener has closed to all listening parties
close(socket.closed)
})
return nil
}
// Addr returns the listener's network address.
func (socket *fakeSocket) Addr() net.Addr {
return &utils.NetAddr{AddrNetwork: "tcp", Addr: "socket.over.ssh"}
}
-7
View File
@@ -258,13 +258,6 @@ func (c *CircularBuffer) removeWatcherWithLock(watcher *BufferWatcher) {
}
}
func max(a, b int) int {
if a > b {
return b
}
return a
}
// BufferWatcher is a watcher connected to the
// buffer and receiving fan-out events from the watcher
type BufferWatcher struct {
-14
View File
@@ -633,20 +633,6 @@ func (b *DynamoDBBackend) createTable(ctx context.Context, tableName string, ran
return trace.Wrap(err)
}
// deleteTable deletes DynamoDB table with a given name
func (b *DynamoDBBackend) deleteTable(ctx context.Context, tableName string, wait bool) error {
tn := aws.String(tableName)
_, err := b.svc.DeleteTable(&dynamodb.DeleteTableInput{TableName: tn})
if err != nil {
return trace.Wrap(err)
}
if wait {
return trace.Wrap(
b.svc.WaitUntilTableNotExists(&dynamodb.DescribeTableInput{TableName: tn}))
}
return nil
}
type getResult struct {
records []record
// lastEvaluatedKey is the primary key of the item where the operation stopped, inclusive of the
-1
View File
@@ -129,7 +129,6 @@ type EtcdBackend struct {
nodes []string
*log.Entry
cfg *Config
etcdKey string
client *clientv3.Client
cancelC chan bool
stopC chan bool
-2
View File
@@ -28,7 +28,6 @@ import (
"github.com/gravitational/teleport/lib/utils"
"github.com/gravitational/trace"
"go.etcd.io/etcd/clientv3"
"gopkg.in/check.v1"
)
@@ -37,7 +36,6 @@ func TestEtcd(t *testing.T) { check.TestingT(t) }
type EtcdSuite struct {
bk *EtcdBackend
suite test.BackendSuite
client *clientv3.Client
config backend.Params
skip bool
}
-5
View File
@@ -672,11 +672,6 @@ func (b *FirestoreBackend) getIndexParent() string {
return "projects/" + b.ProjectID + "/databases/(default)/collectionGroups/" + b.CollectionName
}
// deleteAllItems deletes all items from the database, used in tests
func (b *FirestoreBackend) deleteAllItems() {
DeleteAllDocuments(b.clientContext, b.svc, b.CollectionName)
}
// DeleteAllDocuments will delete all documents in a collection.
func DeleteAllDocuments(ctx context.Context, svc *firestore.Client, collectionName string) {
docs, _ := svc.Collection(collectionName).Documents(ctx).GetAll()
-15
View File
@@ -185,9 +185,6 @@ type LiteBackend struct {
Config
*log.Entry
db *sql.DB
// tx stores the current transaction if this backend instance
// is bound to transaction
tx *sql.Tx
// clock is used to generate time,
// could be swapped in tests for fixed time
clock clockwork.Clock
@@ -878,10 +875,6 @@ func expires(t time.Time) interface{} {
return t.UTC()
}
func nop() error {
return nil
}
func convertError(err error) error {
origError := trace.Unwrap(err)
if isClosedError(origError) {
@@ -918,14 +911,6 @@ func isInterrupt(err error) bool {
return e.Code == sqlite3.ErrInterrupt
}
func isRetryError(err error) bool {
e, ok := trace.Unwrap(err).(sqlite3.Error)
if !ok {
return false
}
return e.Code == sqlite3.ErrBusy || e.Code == sqlite3.ErrLocked
}
func isReadonlyError(err error) bool {
e, ok := trace.Unwrap(err).(sqlite3.Error)
if !ok {
-8
View File
@@ -740,11 +740,3 @@ func ExpectItems(c *check.C, items, expected []backend.Item) {
c.Assert(string(items[i].Value), check.Equals, string(expected[i].Value))
}
}
func toSet(vals []string) map[string]struct{} {
out := make(map[string]struct{}, len(vals))
for _, v := range vals {
out[v] = struct{}{}
}
return out
}
-21
View File
@@ -383,27 +383,6 @@ func (proxy *ProxyClient) ConnectToCluster(ctx context.Context, clusterName stri
return clt, nil
}
// closerConn wraps connection and attaches additional closers to it
type closerConn struct {
net.Conn
closers []io.Closer
}
// addCloser adds any closer in ctx that will be called
// whenever server closes session channel
func (c *closerConn) addCloser(closer io.Closer) {
c.closers = append(c.closers, closer)
}
func (c *closerConn) Close() error {
var errors []error
for _, closer := range c.closers {
errors = append(errors, closer.Close())
}
errors = append(errors, c.Conn.Close())
return trace.NewAggregate(errors...)
}
// nodeName removes the port number from the hostname, if present
func nodeName(node string) string {
n, _, err := net.SplitHostPort(node)
-20
View File
@@ -57,9 +57,6 @@ const (
// SymlinkFilename is a name of the symlink pointing to the last
// current log file
SymlinkFilename = "events.log"
// sessionsMigratedEvent is a sessions migration event used internally
sessionsMigratedEvent = "sessions.migrated"
)
var (
@@ -110,10 +107,6 @@ type AuditLog struct {
// playbackDir is a directory used for unpacked session recordings
playbackDir string
// fileTime is a rounded (to a day, by default) timestamp of the
// currently opened file
fileTime time.Time
// activeDownloads helps to serialize simultaneous downloads
// from the session record server
activeDownloads map[string]context.Context
@@ -932,19 +925,6 @@ func (l *AuditLog) EmitAuditEvent(event Event, fields EventFields) error {
return nil
}
// emitEvent emits event for test purposes
func (l *AuditLog) emitEvent(e AuditLogEvent) {
if l.EventsC == nil {
return
}
select {
case l.EventsC <- &e:
return
default:
l.Warningf("Blocked on the events channel.")
}
}
// auditDirs returns directories used for audit log storage
func (l *AuditLog) auditDirs() ([]string, error) {
authServers, err := l.getAuthServers()
-72
View File
@@ -115,11 +115,6 @@ type event struct {
EventNamespace string
}
type keyLookup struct {
HashKey string
FullPath string
}
const (
// keyExpires is a key used for TTL specification
keyExpires = "Expires"
@@ -133,9 +128,6 @@ const (
// keyEventNamespace
keyEventNamespace = "EventNamespace"
// sectionDefault
sectionDefault = "default"
// keyCreatedAt identifies created at key
keyCreatedAt = "CreatedAt"
@@ -585,53 +577,6 @@ func (b *Log) createTable(tableName string) error {
return trace.Wrap(err)
}
// deleteAllItems deletes all items from the database, used in tests
func (b *Log) deleteAllItems() error {
out, err := b.svc.Scan(&dynamodb.ScanInput{TableName: aws.String(b.Tablename)})
if err != nil {
return trace.Wrap(err)
}
var requests []*dynamodb.WriteRequest
for _, item := range out.Items {
requests = append(requests, &dynamodb.WriteRequest{
DeleteRequest: &dynamodb.DeleteRequest{
Key: map[string]*dynamodb.AttributeValue{
keySessionID: item[keySessionID],
keyEventIndex: item[keyEventIndex],
},
},
})
}
if len(requests) == 0 {
return nil
}
req, _ := b.svc.BatchWriteItemRequest(&dynamodb.BatchWriteItemInput{
RequestItems: map[string][]*dynamodb.WriteRequest{
b.Tablename: requests,
},
})
err = req.Send()
err = convertError(err)
if err != nil {
return trace.Wrap(err)
}
return nil
}
// deleteTable deletes DynamoDB table with a given name
func (b *Log) deleteTable(tableName string, wait bool) error {
tn := aws.String(tableName)
_, err := b.svc.DeleteTable(&dynamodb.DeleteTableInput{TableName: tn})
if err != nil {
return trace.Wrap(err)
}
if wait {
return trace.Wrap(
b.svc.WaitUntilTableNotExists(&dynamodb.DescribeTableInput{TableName: tn}))
}
return nil
}
// Close the DynamoDB driver
func (b *Log) Close() error {
return nil
@@ -660,20 +605,3 @@ func convertError(err error) error {
return err
}
}
type eventlist []event
// Len is part of sort.Interface.
func (e eventlist) Len() int {
return len(e)
}
// Swap is part of sort.Interface.
func (e eventlist) Swap(i, j int) {
e[i], e[j] = e[j], e[i]
}
// Less is part of sort.Interface.
func (e eventlist) Less(i, j int) bool {
return e[i].EventIndex < e[j].EventIndex
}
@@ -29,7 +29,6 @@ import (
func TestFile(t *testing.T) { check.TestingT(t) }
type FileSuite struct {
handler *Handler
test.HandlerSuite
}
@@ -36,7 +36,7 @@ import (
"cloud.google.com/go/firestore"
"cloud.google.com/go/firestore/apiv1/admin"
apiv1 "cloud.google.com/go/firestore/apiv1/admin"
"github.com/gravitational/trace"
"github.com/jonboulle/clockwork"
@@ -523,11 +523,6 @@ func (l *Log) ensureIndexes(adminSvc *apiv1.FirestoreAdminClient) error {
return err
}
// deleteAllItems deletes all items from the database, used in tests
func (l *Log) deleteAllItems() {
firestorebk.DeleteAllDocuments(l.svcContext, l.svc, l.CollectionName)
}
// Close the Firestore driver
func (l *Log) Close() error {
l.svcCancel()
-22
View File
@@ -26,7 +26,6 @@ import (
"time"
"github.com/prometheus/client_golang/prometheus"
"google.golang.org/api/iterator"
"google.golang.org/grpc"
"github.com/gravitational/teleport"
@@ -248,27 +247,6 @@ func (h *Handler) Download(ctx context.Context, sessionID session.ID, writer io.
return nil
}
// delete bucket deletes bucket and all it's contents and is used in tests
// this app should not have the authority to create/destroy resources
func (h *Handler) deleteBucket() error {
objectsIterator := h.gcsClient.Bucket(h.Config.Bucket).Objects(h.clientContext, nil)
for {
attrs, err := objectsIterator.Next()
if err == iterator.Done {
break
}
if err != nil {
return convertGCSError(err)
}
err = h.gcsClient.Bucket(h.Config.Bucket).Object(attrs.Name).Delete(h.clientContext)
if err != nil {
return convertGCSError(err)
}
}
err := h.gcsClient.Bucket(h.Config.Bucket).Delete(h.clientContext)
return convertGCSError(err)
}
func (h *Handler) path(sessionID session.ID) string {
if h.Path == "" {
return string(sessionID) + ".tar"
-32
View File
@@ -128,7 +128,6 @@ type DiskSessionLogger struct {
sync.Mutex
sid session.ID
sessionDir string
indexFile *os.File
@@ -633,37 +632,6 @@ func newGzipWriter(file *os.File) *gzipWriter {
}
}
// gzipReader wraps file, on close close both gzip writer and file
type gzipReader struct {
io.ReadCloser
file io.Closer
}
// Close closes file and gzip writer
func (f *gzipReader) Close() error {
var errors []error
if f.ReadCloser != nil {
errors = append(errors, f.ReadCloser.Close())
f.ReadCloser = nil
}
if f.file != nil {
errors = append(errors, f.file.Close())
f.file = nil
}
return trace.NewAggregate(errors...)
}
func newGzipReader(file *os.File) (*gzipReader, error) {
reader, err := gzip.NewReader(file)
if err != nil {
return nil, trace.Wrap(err)
}
return &gzipReader{
ReadCloser: reader,
file: file,
}, nil
}
const (
// eventsSuffix is the suffix of the archive that contians session events.
eventsSuffix = "events.gz"
-10
View File
@@ -270,16 +270,6 @@ func (h *portForwardProxy) monitorStreamPair(p *httpStreamPair, timeout <-chan t
h.removeStreamPair(p.requestID)
}
// hasStreamPair returns a bool indicating if a stream pair for requestID
// exists.
func (h *portForwardProxy) hasStreamPair(requestID string) bool {
h.streamPairsLock.RLock()
defer h.streamPairsLock.RUnlock()
_, ok := h.streamPairs[requestID]
return ok
}
// removeStreamPair removes the stream pair identified by requestID from streamPairs.
func (h *portForwardProxy) removeStreamPair(requestID string) {
h.streamPairsLock.Lock()
-15
View File
@@ -21,7 +21,6 @@ import (
"bytes"
"context"
"crypto/tls"
"encoding/base64"
"fmt"
"io"
"io/ioutil"
@@ -67,10 +66,6 @@ type SpdyRoundTripper struct {
// dialWithContext is the function used connect to remote address
dialWithContext func(context context.Context, network, address string) (net.Conn, error)
// proxier knows which proxy to use given a request, defaults to http.ProxyFromEnvironment
// Used primarily for mocking the proxy discovery in tests.
proxier func(req *http.Request) (*url.URL, error)
// followRedirects indicates if the round tripper should examine responses for redirects and
// follow them.
followRedirects bool
@@ -172,16 +167,6 @@ func (s *SpdyRoundTripper) dialWithoutProxy(url *url.URL) (net.Conn, error) {
return conn, nil
}
// proxyAuth returns, for a given proxy URL, the value to be used for the Proxy-Authorization header
func (s *SpdyRoundTripper) proxyAuth(proxyURL *url.URL) string {
if proxyURL == nil || proxyURL.User == nil {
return ""
}
credentials := proxyURL.User.String()
encodedAuth := base64.StdEncoding.EncodeToString([]byte(credentials))
return fmt.Sprintf("Basic %s", encodedAuth)
}
// RoundTrip executes the Request and upgrades it. After a successful upgrade,
// clients may call SpdyRoundTripper.Connection() to retrieve the upgraded
// connection.
-21
View File
@@ -190,12 +190,6 @@ func (a *Agent) String() string {
return fmt.Sprintf("agent(id=%d,state=%v) -> %v:%v, discover %v", a.ID, a.getState(), a.ClusterName, a.Addr.String(), Proxies(a.DiscoverProxies))
}
func (a *Agent) getLastStateChange() time.Time {
a.RLock()
defer a.RUnlock()
return a.stateChange
}
func (a *Agent) setStateAndPrincipals(state string, principals []string) {
a.Lock()
defer a.Unlock()
@@ -236,21 +230,6 @@ func (a *Agent) Wait() error {
return nil
}
func (a *Agent) isDiscovering(proxy services.Server) bool {
for _, discoverProxy := range a.DiscoverProxies {
if a.getState() != agentStateDiscovering && a.getState() != agentStateConnecting {
continue
}
proxyID := fmt.Sprintf("%v.%v", proxy.GetName(), a.ClusterName)
discoverID := fmt.Sprintf("%v.%v", discoverProxy.GetName(), a.ClusterName)
if proxyID == discoverID {
return true
}
}
return false
}
// connectedTo returns true if connected services.Server passed in.
func (a *Agent) connectedTo(proxy services.Server) bool {
principals := a.getPrincipals()
-9
View File
@@ -186,15 +186,6 @@ func (m *AgentPool) processSeekEvents() {
}
}
func foundInOneOf(proxy services.Server, agents []*Agent) bool {
for _, agent := range agents {
if agent.isDiscovering(proxy) || agent.connectedTo(proxy) {
return true
}
}
return false
}
// FetchAndSyncAgents executes one time fetch and sync request
// (used in tests instead of polling)
func (m *AgentPool) FetchAndSyncAgents() error {
-14
View File
@@ -26,7 +26,6 @@ import (
"golang.org/x/crypto/ssh"
"github.com/gravitational/teleport"
"github.com/gravitational/teleport/lib/auth"
"github.com/gravitational/teleport/lib/services"
"github.com/gravitational/teleport/lib/utils"
@@ -176,14 +175,6 @@ func (c *remoteConn) setLastHeartbeat(tm time.Time) {
atomic.StoreInt64(&c.lastHeartbeat, tm.UnixNano())
}
func (c *remoteConn) getLastHeartbeat() time.Time {
unixNano := atomic.LoadInt64(&c.lastHeartbeat)
if unixNano == 0 {
return time.Time{}
}
return time.Unix(0, unixNano)
}
// isReady returns true when connection is ready to be tried,
// it returns true when connection has received the first heartbeat
func (c *remoteConn) isReady() bool {
@@ -257,8 +248,3 @@ func (c *remoteConn) sendDiscoveryRequest(req discoveryRequest) error {
return nil
}
func (c *remoteConn) isOnline(conn services.TunnelConnection) bool {
tunnelStatus := services.TunnelConnectionStatus(c.clock, conn, c.offlineThreshold)
return tunnelStatus == teleport.RemoteClusterStatusOnline
}
-4
View File
@@ -74,10 +74,6 @@ func (proxies Proxies) Equal(other []services.Server) bool {
return true
}
func (r discoveryRequest) key() agentKey {
return agentKey{clusterName: r.ClusterName, tunnelType: r.Type, addr: r.ClusterAddr}
}
func (r discoveryRequest) String() string {
return fmt.Sprintf("discovery request, cluster name: %v, address: %v, proxies: %v",
r.ClusterName, r.ClusterAddr, Proxies(r.Proxies))
+3 -10
View File
@@ -17,7 +17,6 @@ limitations under the License.
package reversetunnel
import (
"context"
"fmt"
"net"
"sync"
@@ -77,12 +76,9 @@ func newlocalSite(srv *server, domainName string, client auth.ClientI) (*localSi
type localSite struct {
sync.Mutex
authServer string
log *log.Entry
domainName string
connections []*remoteConn
lastUsed int
srv *server
log *log.Entry
domainName string
srv *server
// client provides access to the Auth Server API of the local cluster.
client auth.ClientI
@@ -96,9 +92,6 @@ type localSite struct {
// remoteConns maps UUID to a remote connection.
remoteConns map[string]*remoteConn
// closeContext is used to signal when the site is shutting down.
closeContext context.Context
// clock is used to control time in tests.
clock clockwork.Clock
+6 -7
View File
@@ -239,13 +239,12 @@ func (s *GroupHandle) Gossip() chan<- string {
// proxyGroup manages all proxy seekers for a group.
type proxyGroup struct {
sync.Mutex
id Key
conf Config
states map[string]seeker
proxyC <-chan string
seekC chan<- Key
statC chan Status
lastStat *Status
id Key
conf Config
states map[string]seeker
proxyC <-chan string
seekC chan<- Key
statC chan Status
}
// run is the "main loop" for the seek process.
+2 -8
View File
@@ -19,11 +19,12 @@ package seek
import (
"context"
"fmt"
"gopkg.in/check.v1"
pr "math/rand"
"sync"
"testing"
"time"
"gopkg.in/check.v1"
)
type simpleTestProxies struct {
@@ -113,13 +114,6 @@ func newTestProxy(life time.Duration) testProxy {
return testProxy{principals, life}
}
// prMillis generates a pseudorandom duration
// of [0,max) milliseconds.
func prMillis(min int64, max int64) time.Duration {
n := time.Duration(pr.Int63n(max-min) + min)
return n * time.Millisecond
}
func prDuration(min time.Duration, max time.Duration) time.Duration {
mn, mx := int64(min), int64(max)
rslt := pr.Int63n(mx-mn) + mn
-80
View File
@@ -655,38 +655,6 @@ func (s *server) findLocalCluster(sconn *ssh.ServerConn) (*localSite, error) {
return nil, trace.BadParameter("local cluster %v not found", clusterName)
}
// isHostAuthority is called during checking the client key, to see if the signing
// key is the real host CA authority key.
func (s *server) isHostAuthority(auth ssh.PublicKey, address string) bool {
keys, err := s.getTrustedCAKeys(services.HostCA)
if err != nil {
s.Errorf("failed to retrieve trusted keys, err: %v", err)
return false
}
for _, k := range keys {
if sshutils.KeysEqual(k, auth) {
return true
}
}
return false
}
// isUserAuthority is called during checking the client key, to see if the signing
// key is the real user CA authority key.
func (s *server) isUserAuthority(auth ssh.PublicKey) bool {
keys, err := s.getTrustedCAKeys(services.UserCA)
if err != nil {
s.Errorf("failed to retrieve trusted keys, err: %v", err)
return false
}
for _, k := range keys {
if sshutils.KeysEqual(k, auth) {
return true
}
}
return false
}
func (s *server) getTrustedCAKeysByID(id services.CertAuthID) ([]ssh.PublicKey, error) {
ca, err := s.localAccessPoint.GetCertAuthority(id, false, services.SkipValidation())
if err != nil {
@@ -695,22 +663,6 @@ func (s *server) getTrustedCAKeysByID(id services.CertAuthID) ([]ssh.PublicKey,
return ca.Checkers()
}
func (s *server) getTrustedCAKeys(certType services.CertAuthType) ([]ssh.PublicKey, error) {
cas, err := s.localAccessPoint.GetCertAuthorities(certType, false, services.SkipValidation())
if err != nil {
return nil, err
}
out := []ssh.PublicKey{}
for _, ca := range cas {
checkers, err := ca.Checkers()
if err != nil {
return nil, trace.Wrap(err)
}
out = append(out, checkers...)
}
return out, nil
}
func (s *server) keyAuth(conn ssh.ConnMetadata, key ssh.PublicKey) (*ssh.Permissions, error) {
logger := s.WithFields(log.Fields{
"remote": conn.RemoteAddr(),
@@ -1011,38 +963,6 @@ func newRemoteSite(srv *server, domainName string, sconn ssh.Conn) (*remoteSite,
return remoteSite, nil
}
// sendVersionRequest sends a request for the version remote Teleport cluster.
// If a response is not received within one second, it's assumed it's a legacy
// cluster.
func sendVersionRequest(sconn ssh.Conn, closeContext context.Context) (string, error) {
errorCh := make(chan error, 1)
versionCh := make(chan string, 1)
go func() {
ok, payload, err := sconn.SendRequest(versionRequest, true, nil)
if err != nil {
errorCh <- err
return
}
if !ok {
errorCh <- trace.BadParameter("no response to %v request", versionRequest)
return
}
versionCh <- string(payload)
}()
select {
case ver := <-versionCh:
return ver, nil
case err := <-errorCh:
return "", trace.Wrap(err)
case <-time.After(defaults.ReadHeadersTimeout):
return "", trace.BadParameter("timeout waiting for version")
case <-closeContext.Done():
return "", closeContext.Err()
}
}
const (
extHost = "host@teleport"
extCertType = "certtype@teleport"
-6
View File
@@ -2360,12 +2360,6 @@ func (process *TeleportProcess) setReporter(reporter *backend.Reporter) {
process.reporter = reporter
}
func (process *TeleportProcess) getReporter() *backend.Reporter {
process.Lock()
defer process.Unlock()
return process.reporter
}
// WaitWithContext waits until all internal services stop.
func (process *TeleportProcess) WaitWithContext(ctx context.Context) {
local, cancel := context.WithCancel(ctx)
+1 -2
View File
@@ -119,7 +119,6 @@ type LocalSupervisor struct {
sync.Mutex
wg *sync.WaitGroup
services []Service
errors []error
events map[string]Event
eventsC chan Event
eventWaiters map[string][]*waiter
@@ -467,5 +466,5 @@ type ServiceFunc func() error
const (
stateCreated = iota
stateStarted = iota
stateStarted
)
-70
View File
@@ -1,70 +0,0 @@
/*
Copyright 2015 Gravitational, Inc.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
You may obtain a copy of the License at
http://www.apache.org/licenses/LICENSE-2.0
Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.
*/
package service
import (
"bytes"
"fmt"
"io/ioutil"
"os"
"strings"
"text/template"
"github.com/gravitational/trace"
)
func renderTemplate(data []byte) ([]byte, error) {
t, err := template.New("tpl").Parse(string(data))
if err != nil {
return nil, trace.Wrap(err)
}
buf := &bytes.Buffer{}
c, err := newCtx()
if err != nil {
return nil, trace.Wrap(err)
}
if err := t.Execute(buf, c); err != nil {
return nil, trace.Wrap(err)
}
return buf.Bytes(), nil
}
func newCtx() (*ctx, error) {
values := os.Environ()
c := &ctx{
env: make(map[string]string, len(values)),
}
for _, v := range values {
vals := strings.SplitN(v, "=", 2)
if len(vals) != 2 {
return nil, trace.Errorf("failed to parse variable: '%v'", v)
}
c.env[vals[0]] = vals[1]
}
return c, nil
}
type ctx struct {
env map[string]string
}
func (c *ctx) File(path string) (string, error) {
o, err := ioutil.ReadFile(path)
if err != nil {
return "", trace.Wrap(err, fmt.Sprintf("reading file: %v", path))
}
return string(o), nil
}
-3
View File
@@ -40,9 +40,6 @@ type CertAuthorityV2 struct {
Metadata Metadata `json:"metadata"`
// Spec contains cert authority specification
Spec CertAuthoritySpecV2 `json:"spec"`
// rawObject is object that is raw object stored in DB
// without any conversions applied, used in migrations
rawObject interface{}
}
// Rotation is a status of the rotation of the certificate authority
-36
View File
@@ -17,11 +17,9 @@ limitations under the License.
package services
import (
"bytes"
"encoding/json"
"fmt"
"net/url"
"text/template"
"time"
"github.com/gravitational/teleport"
@@ -491,40 +489,6 @@ func (o *OIDCConnectorV2) MapClaims(claims jose.Claims) []string {
return utils.Deduplicate(roles)
}
func executeStringTemplate(raw string, claims jose.Claims) (string, error) {
tmpl, err := template.New("dynamic-roles").Parse(raw)
if err != nil {
return "", trace.Wrap(err)
}
var buf bytes.Buffer
err = tmpl.Execute(&buf, claims)
if err != nil {
return "", trace.Wrap(err)
}
return buf.String(), nil
}
func executeSliceTemplate(raw []string, claims jose.Claims) ([]string, error) {
var sl []string
for _, v := range raw {
tmpl, err := template.New("dynamic-roles").Parse(v)
if err != nil {
return nil, trace.Wrap(err)
}
var buf bytes.Buffer
err = tmpl.Execute(&buf, claims)
if err != nil {
return nil, trace.Wrap(err)
}
sl = append(sl, buf.String())
}
return sl, nil
}
// Check returns nil if all parameters are great, err otherwise
func (o *OIDCConnectorV2) Check() error {
if o.Metadata.Name == "" {
-17
View File
@@ -800,20 +800,3 @@ type SigningKeyPair struct {
// Cert is certificate in OpenSSH authorized keys format
Cert string `json:"cert"`
}
// buildAssertionMap takes an saml2.AssertionInfo and builds a friendly map
// that can be used to access assertion/value pairs. If multiple values are
// returned for an assertion, they are joined into a string by ",".
func buildAssertionMap(assertionInfo saml2.AssertionInfo) map[string]string {
assertionMap := make(map[string]string)
for _, assr := range assertionInfo.Values {
var vals []string
for _, v := range assr.Values {
vals = append(vals, v.Value)
}
assertionMap[assr.Name] = strings.Join(vals, ",")
}
return assertionMap
}
-20
View File
@@ -113,26 +113,6 @@ func (s *ServicesTestSuite) Users() services.UsersService {
return s.UsersS
}
func (s *ServicesTestSuite) collectChanges(c *check.C, expected int) []interface{} {
changes := make([]interface{}, expected)
for i := range changes {
select {
case changes[i] = <-s.ChangesC:
// successfully collected changes
case <-time.After(2 * time.Second):
c.Fatalf("Timeout occurred waiting for events")
}
}
return changes
}
func (s *ServicesTestSuite) expectChanges(c *check.C, expected ...interface{}) {
changes := s.collectChanges(c, len(expected))
for i, ch := range changes {
c.Assert(ch, check.DeepEquals, expected[i])
}
}
func userSlicesEqual(c *check.C, a []services.User, b []services.User) {
comment := check.Commentf("a: %#v b: %#v", a, b)
c.Assert(len(a), check.Equals, len(b), comment)
-8
View File
@@ -428,14 +428,6 @@ func (h *AuthHandlers) isProxy() bool {
return false
}
func (h *AuthHandlers) isTeleportNode() bool {
if h.Component == teleport.ComponentNode {
return true
}
return false
}
// extractRolesFromCert extracts roles from certificate metadata extensions.
func extractRolesFromCert(cert *ssh.Certificate) ([]string, error) {
data, ok := cert.Extensions[teleport.CertExtensionTeleportRoles]
-6
View File
@@ -765,12 +765,6 @@ func closeAll(closers ...io.Closer) error {
return trace.NewAggregate(errs...)
}
type closerFunc func() error
func (f closerFunc) Close() error {
return f()
}
// NewTrackingReader returns a new instance of
// activity tracking reader.
func NewTrackingReader(ctx *ServerContext, r io.Reader) *TrackingReader {
-5
View File
@@ -33,7 +33,6 @@ import (
"golang.org/x/crypto/ssh"
"github.com/gravitational/teleport"
"github.com/gravitational/teleport/lib/bpf"
"github.com/gravitational/teleport/lib/events"
"github.com/gravitational/teleport/lib/services"
"github.com/gravitational/teleport/lib/utils"
@@ -121,10 +120,6 @@ type localExec struct {
// Ctx holds the *ServerContext.
Ctx *ServerContext
// sessionContext holds the BPF session context used to lookup and interact
// with BPF sessions.
sessionContext *bpf.SessionContext
}
// GetCommand returns the command string.
-13
View File
@@ -24,7 +24,6 @@ import (
"os"
os_exec "os/exec"
"os/user"
"path"
"path/filepath"
"strconv"
"testing"
@@ -325,18 +324,6 @@ func (s *ExecSuite) OpenChannel(string, []byte) (ssh.Channel, <-chan *ssh.Reques
}
func (s *ExecSuite) Wait() error { return nil }
// findExecutable helper finds a given executable name (like 'ls') in $PATH
// and returns the full path
func findExecutable(execName string) string {
for _, dir := range filepath.SplitList(os.Getenv("PATH")) {
fp := path.Join(dir, execName)
if utils.IsFile(fp) {
return fp
}
}
return "not found in $PATH: " + execName
}
type fakeTerminal struct {
f *os.File
}
-1
View File
@@ -258,7 +258,6 @@ func newFakeAnnouncer(ctx context.Context) *fakeAnnouncer {
type fakeAnnouncer struct {
err error
srv services.Server
upsertCalls map[HeartbeatMode]int
closeCalls int
ctx context.Context
-9
View File
@@ -44,11 +44,6 @@ const (
// in a terminal) to be instanly replayed to the newly joining
// parties
instantReplayLen = 20
// maxTermSyncErrorCount defines how many subsequent erorrs
// we should tolerate before giving up trying to sync the
// term size
maxTermSyncErrorCount = 5
)
var (
@@ -461,10 +456,6 @@ type session struct {
// client hits "page refresh").
lingerTTL time.Duration
// termSizeC is used to push terminal resize events from SSH "on-size-changed"
// event handler into "push-to-web-client" loop.
termSizeC chan []byte
// login stores the login of the initial session creator
login string
+2 -3
View File
@@ -99,9 +99,8 @@ func makeFileInfo(filePath string) (FileInfo, error) {
// localFileInfo is implementation of FileInfo for local files
type localFileInfo struct {
isRecursive bool
filePath string
fileInfo os.FileInfo
filePath string
fileInfo os.FileInfo
}
// IsDir tells this is a directory
-19
View File
@@ -504,25 +504,6 @@ func (cmd *command) receiveDir(st *state, fc newFileCmd, ch io.ReadWriter) error
return nil
}
// sendError gets called during all errors during SCP transmission.
// It writes it back to the SCP client
func (cmd *command) sendError(ch io.ReadWriter, err error) error {
if err == nil {
return nil
}
cmd.log.Error(err)
message := err.Error()
bytes := make([]byte, 0, len(message)+2)
bytes = append(bytes, ErrByte)
bytes = append(bytes, message...)
bytes = append(bytes, []byte{'\n'}...)
_, writeErr := ch.Write(bytes)
if writeErr != nil {
cmd.log.Error(writeErr)
}
return trace.Wrap(err)
}
type newFileCmd struct {
Mode int64
Length uint64
-4
View File
@@ -218,10 +218,6 @@ func urlToNetAddr(u string) NetAddr {
return *MustParseAddr(parsed.Host)
}
func localURL(port string) string {
return fmt.Sprintf("http://127.0.0.1:%v", port)
}
func localAddr(port string) NetAddr {
return *MustParseAddr(fmt.Sprintf("127.0.0.1:%v", port))
}
+1 -2
View File
@@ -31,8 +31,7 @@ import (
// TimeoutSuite helps us to test ObeyTimeout mechanism. We use HTTP server/client
// machinery to test timeouts
type TimeoutSuite struct {
lastRequest *http.Request
server *httptest.Server
server *httptest.Server
}
var _ = check.Suite(&TimeoutSuite{})
-15
View File
@@ -55,7 +55,6 @@ import (
"github.com/jonboulle/clockwork"
"github.com/julienschmidt/httprouter"
lemma_secret "github.com/mailgun/lemma/secret"
"github.com/mailgun/ttlmap"
log "github.com/sirupsen/logrus"
"github.com/tstranex/u2f"
"golang.org/x/crypto/ssh"
@@ -67,7 +66,6 @@ type Handler struct {
httprouter.Router
cfg Config
auth *sessionCache
sites *ttlmap.TtlMap
sessionStreamPollPeriod time.Duration
clock clockwork.Clock
}
@@ -1413,11 +1411,6 @@ func (h *Handler) getSiteNamespaces(w http.ResponseWriter, r *http.Request, _ ht
}, nil
}
type nodeWithSessions struct {
Node services.ServerV1 `json:"node"`
Sessions []session.Session `json:"sessions"`
}
func (h *Handler) siteNodesGet(w http.ResponseWriter, r *http.Request, p httprouter.Params, ctx *SessionContext, site reversetunnel.RemoteSite) (interface{}, error) {
namespace := p.ByName("namespace")
if !services.IsValidNamespace(namespace) {
@@ -1539,10 +1532,6 @@ func (h *Handler) siteSessionGenerate(w http.ResponseWriter, r *http.Request, p
return siteSessionGenerateResponse{Session: req.Session}, nil
}
type siteSessionUpdateReq struct {
TerminalParams session.TerminalParams `json:"terminal_params"`
}
type siteSessionsGetResponse struct {
Sessions []session.Session `json:"sessions"`
}
@@ -1726,10 +1715,6 @@ func queryLimit(query url.Values, name string, def int) (int, error) {
return limit, nil
}
type siteSessionStreamGetResponse struct {
Bytes []byte `json:"bytes"`
}
// siteSessionStreamGet returns a byte array from a session's stream
//
// GET /v1/webapi/sites/:site/namespaces/:namespace/sessions/:sid/stream?query
-2
View File
@@ -68,7 +68,6 @@ import (
"github.com/gravitational/trace"
"github.com/beevik/etree"
"github.com/gokyle/hotp"
"github.com/golang/protobuf/proto"
"github.com/jonboulle/clockwork"
lemma_secret "github.com/mailgun/lemma/secret"
@@ -270,7 +269,6 @@ type authPack struct {
otpSecret string
user string
login string
otp *hotp.HOTP
session *CreateSessionResponse
clt *client.WebClient
cookies []*http.Cookie
-3
View File
@@ -152,9 +152,6 @@ type TerminalHandler struct {
// sshSession holds the "shell" SSH channel to the node.
sshSession *ssh.Session
// teleportClient is the client used to form the connection.
teleportClient *client.TeleportClient
// terminalContext is used to signal when the terminal sesson is closing.
terminalContext context.Context