attempt to fix deadlock by releasing view lock before broadcasting CreateShard

This commit is contained in:
Matt Jaffee 2018-12-11 15:00:46 -06:00
parent bb65a4a14f
commit 4b786e1057
No known key found for this signature in database
GPG key ID: 08A3DFFF987B11BF
2 changed files with 31 additions and 14 deletions

View file

@ -1504,7 +1504,15 @@ func GetHTTPClient(t *tls.Config) *http.Client {
if t != nil {
transport.TLSClientConfig = t
}
return &http.Client{Transport: transport}
return &http.Client{
Transport: transport,
// Internal queries will time out after 2h 7m by default. This is
// reduced from the old default of no timeout, so it was thought we
// should keep it fairly high, but it could probably be reduced in most
// cases. It is set to an odd number in the hopes that it will be
// recognizable in stats/traces/logs when this limit is being hit.
Timeout: 127 * time.Minute,
}
}
// handlPostRoaringImport

35
view.go
View file

@ -206,37 +206,46 @@ func (v *view) recalculateCaches() {
// CreateFragmentIfNotExists returns a fragment in the view by shard.
func (v *view) CreateFragmentIfNotExists(shard uint64) (*fragment, error) {
v.mu.Lock()
defer v.mu.Unlock()
return v.createFragmentIfNotExists(shard)
frag, msg, err := v.createFragmentIfNotExists(shard)
if err == nil && msg != nil {
// Broadcast a message that a new max shard was just created.
if err = v.broadcaster.SendSync(msg); err != nil {
v.mu.Lock()
delete(v.fragments, shard)
v.mu.Unlock()
frag.close()
return nil, errors.Wrap(err, "sending createshard message")
}
}
return frag, err
}
func (v *view) createFragmentIfNotExists(shard uint64) (*fragment, error) {
func (v *view) createFragmentIfNotExists(shard uint64) (*fragment, *CreateShardMessage, error) {
v.mu.Lock()
defer v.mu.Unlock()
// Find fragment in cache first.
if frag := v.fragments[shard]; frag != nil {
return frag, nil
return frag, nil, nil
}
// Initialize and open fragment.
frag := v.newFragment(v.fragmentPath(shard), shard)
if err := frag.Open(); err != nil {
return nil, errors.Wrap(err, "opening fragment")
return nil, nil, errors.Wrap(err, "opening fragment")
}
frag.RowAttrStore = v.rowAttrStore
// Broadcast a message that a new max shard was just created.
if err := v.broadcaster.SendSync(&CreateShardMessage{
msg := &CreateShardMessage{
Index: v.index,
Field: v.field,
Shard: shard,
}); err != nil {
frag.close()
return nil, errors.Wrap(err, "sending createshard message")
}
v.fragments[shard] = frag
// Save to lookup.
v.fragments[shard] = frag
return frag, nil
return frag, msg, nil
}
func (v *view) newFragment(path string, shard uint64) *fragment {