From 74fefb2fdb578fb5dbc760f66593c81ac81dc410 Mon Sep 17 00:00:00 2001 From: jacob Date: Thu, 16 Feb 2023 13:19:35 -0600 Subject: [PATCH] adding debug --- etcd/embed.go | 17 +++++++++++++++-- etcd/leasedkv.go | 3 +++ 2 files changed, 18 insertions(+), 2 deletions(-) diff --git a/etcd/embed.go b/etcd/embed.go index 3eefab6ec..5e0d936cb 100644 --- a/etcd/embed.go +++ b/etcd/embed.go @@ -261,7 +261,7 @@ func (e *Etcd) Close() error { return nil } -const etcdRetryTimes = 3 +const etcdRetryTimes = 4 // newClient requests a new client which is different from the one // passed in. if we've already changed our client (say, because someone @@ -325,7 +325,7 @@ func (e *Etcd) retryClient(fn func(cli *clientv3.Client) error) (err error) { if tries < etcdRetryTimes { retrying = fmt.Sprintf(" (retrying, n=%d)", tries) } - e.logger.Warnf("timeout (%v elapsed) on etcd query%s", elapsed, retrying) + e.logger.Warnf("timeout (%v elapsed) on etcd query %s", elapsed, retrying) // Sleep just a touch longer to give things a time to // stabilize. We're mostly relying on the fact that this is a // timeout to give us a reasonable backoff period and keep us @@ -452,6 +452,7 @@ func (e *Etcd) Start(ctx context.Context) (_ disco.InitialClusterState, err erro // startHeartbeatAndWatcher spins up the heartbeat, and also a background // watcher that watches for changes to events we care about. func (e *Etcd) startHeartbeatAndWatcher(ctx context.Context) error { + e.logger.Debugf("Starting Heartbeat and Watcher...") key := heartbeatPrefix + e.service.ID() e.heartbeatLeasedKV = newLeasedKV(e, e.childContext, key, e.options.HeartbeatTTL) @@ -711,6 +712,7 @@ func (e *Etcd) watchNodes() { // should get cancelled when this Etcd gets shut down. for e.childContext.Err() == nil { err := e.retryClient(func(cli *clientv3.Client) error { + e.logger.Debugf("Calling watchNodesOnce from etcd retryClient...") return e.watchNodesOnce(cli) }) if err != nil { @@ -817,6 +819,7 @@ func (e *Etcd) CreateIndex(ctx context.Context, name string, val []byte) error { // Check for key existence, and execute Op within a transaction. var resp *clientv3.TxnResponse err := e.retryClient(func(cli *clientv3.Client) (err error) { + e.logger.Debugf("CreateIndex call from etcd retryClient...") resp, err = cli.Txn(ctx). If(clientv3util.KeyMissing(key)). Then(op). @@ -843,6 +846,7 @@ func (e *Etcd) DeleteIndex(ctx context.Context, name string) (err error) { // Deleting index and fields in one transaction. err = e.retryClient(func(cli *clientv3.Client) error { + e.logger.Debugf("DeleteIndex call from etcd retryClient...") _, err = cli.Txn(ctx). If(clientv3.Compare(clientv3.Version(key), ">", -1)). Then( @@ -871,6 +875,7 @@ func (e *Etcd) CreateField(ctx context.Context, indexName string, name string, f var resp *clientv3.TxnResponse err := e.retryClient(func(cli *clientv3.Client) (err error) { + e.logger.Debugf("CreateField call from etcd retryClient...") resp, err = cli.Txn(ctx). If( clientv3util.KeyMissing(key)). @@ -900,6 +905,7 @@ func (e *Etcd) UpdateField(ctx context.Context, indexName string, name string, f var resp *clientv3.TxnResponse err := e.retryClient(func(cli *clientv3.Client) (err error) { + e.logger.Debugf("UpdateField call from etcd retryClient...") resp, err = cli.Txn(ctx). If( clientv3util.KeyExists(key)). @@ -924,6 +930,7 @@ func (e *Etcd) DeleteField(ctx context.Context, indexName string, name string) ( // Deleting field and views in one transaction. err = e.retryClient(func(cli *clientv3.Client) (err error) { + e.logger.Debugf("DeleteField call from etcd retryClient...") resp, err = cli.Txn(ctx). If( clientv3.Compare(clientv3.Version(key), ">", -1)). @@ -957,6 +964,7 @@ func (e *Etcd) CreateView(ctx context.Context, indexName, fieldName, name string // Check for key existence, and execute Op within a transaction. err = e.retryClient(func(cli *clientv3.Client) (err error) { + e.logger.Debugf("CreateView call from etcd retryClient...") _, err = cli.Txn(ctx). If(clientv3util.KeyMissing(key)). Then(clientv3.OpPut(key, "")). @@ -976,6 +984,7 @@ func (e *Etcd) DeleteView(ctx context.Context, indexName, fieldName, name string func (e *Etcd) putKey(ctx context.Context, key, val string, opts ...clientv3.OpOption) error { err := e.retryClient(func(cli *clientv3.Client) (err error) { + e.logger.Debugf("putKey call from etcd retryClient...") _, err = cli.Txn(ctx). Then(clientv3.OpPut(key, val, opts...)). Commit() @@ -989,6 +998,7 @@ func (e *Etcd) getKeyBytes(ctx context.Context, key string) ([]byte, error) { op := clientv3.OpGet(key) var resp *clientv3.TxnResponse err := e.retryClient(func(cli *clientv3.Client) (err error) { + e.logger.Debugf("getKeyBytes call from etcd retryClient...") resp, err = cli.Txn(ctx).Then(op).Commit() return err }) @@ -1012,6 +1022,7 @@ func (e *Etcd) getKeyWithPrefix(ctx context.Context, key string) (keys []string, op := clientv3.OpGet(key, clientv3.WithPrefix()) var resp *clientv3.TxnResponse err = e.retryClient(func(cli *clientv3.Client) (err error) { + e.logger.Debugf("getkeyWithPrefix call from etcd retryClient...") resp, err = cli.Txn(ctx).Then(op).Commit() return err }) @@ -1041,6 +1052,7 @@ func (e *Etcd) getKeyWithPrefix(ctx context.Context, key string) (keys []string, func (e *Etcd) keyExists(ctx context.Context, key string) (bool, error) { var resp *clientv3.TxnResponse err := e.retryClient(func(cli *clientv3.Client) (err error) { + e.logger.Debugf("keyExists call from etcd retryClient...") resp, err = cli.Txn(ctx). If(clientv3util.KeyExists(key)). Then(clientv3.OpGet(key, clientv3.WithCountOnly())). @@ -1089,6 +1101,7 @@ func (e *Etcd) SetShards(ctx context.Context, index, field string, shards []byte op := clientv3.OpPut(key, "") op.WithValueBytes(shards) return e.retryClient(func(cli *clientv3.Client) (err error) { + e.logger.Debugf("SetShards call from etcd retryClient...") _, err = cli.Txn(ctx).Then(op).Commit() return }) diff --git a/etcd/leasedkv.go b/etcd/leasedkv.go index 826f4bcc8..298bcff49 100644 --- a/etcd/leasedkv.go +++ b/etcd/leasedkv.go @@ -61,6 +61,7 @@ func (l *leasedKV) Start(initValue string) error { // done channel for the internal context, which will be readable as soon as either // our context or the parent context is done. func (l *leasedKV) create(initValue string) (<-chan *clientv3.LeaseKeepAliveResponse, error) { + log.Printf("trying to create lease and KeepAlive channel...") ctx, cancel := context.WithCancel(l.parentContext) if l.cancel != nil { @@ -104,6 +105,8 @@ func (l *leasedKV) create(initValue string) (<-chan *clientv3.LeaseKeepAliveResp } func (l *leasedKV) consumeLease(ch <-chan *clientv3.LeaseKeepAliveResponse) { + log.Printf("consumeLease from leasedKV %v", l.leaseID) + for { select { case _, ok := <-ch: