PLT-6472: Basic Elastic Search implementation. (#6382)

* PLT-6472: Basic Elastic Search implementation.

This currently supports indexing of posts at create/update/delete time.
It does not support batch indexing or reindexing, and does not support
any entities other than posts yet. The purpose is to more-or-less
replicate the existing full-text search feature but with some of the
immediate benefits of using elastic search.

* Alter settings for AWS compatability.

* Remove unneeded i18n strings.
Этот коммит содержится в:
George Goldberg
2017-05-18 16:26:52 +01:00
коммит произвёл Harrison Healey
родитель 2bbedd9def
Коммит 0db5e3922f
10 изменённых файлов: 343 добавлений и 37 удалений

Просмотреть файл

@@ -125,6 +125,11 @@ func CreatePost(post *model.Post, teamId string, triggerWebhooks bool) (*model.P
rpost = result.Data.(*model.Post) rpost = result.Data.(*model.Post)
} }
esInterface := einterfaces.GetElasticSearchInterface()
if (esInterface != nil && *utils.Cfg.ElasticSearchSettings.EnableIndexing) {
go esInterface.IndexPost(rpost, teamId)
}
if einterfaces.GetMetricsInterface() != nil { if einterfaces.GetMetricsInterface() != nil {
einterfaces.GetMetricsInterface().IncrementPostCreate() einterfaces.GetMetricsInterface().IncrementPostCreate()
} }
@@ -308,6 +313,17 @@ func UpdatePost(post *model.Post, safeUpdate bool) (*model.Post, *model.AppError
} else { } else {
rpost := result.Data.(*model.Post) rpost := result.Data.(*model.Post)
esInterface := einterfaces.GetElasticSearchInterface()
if (esInterface != nil && *utils.Cfg.ElasticSearchSettings.EnableIndexing) {
go func() {
if rchannel := <-Srv.Store.Channel().GetForPost(rpost.Id); rchannel.Err != nil {
l4g.Error("Couldn't get channel %v for post %v for ElasticSearch indexing.", rpost.ChannelId, rpost.Id)
} else {
esInterface.IndexPost(rpost, rchannel.Data.(*model.Channel).TeamId)
}
}()
}
sendUpdatedPostEvent(rpost) sendUpdatedPostEvent(rpost)
InvalidateCacheForChannelPosts(rpost.ChannelId) InvalidateCacheForChannelPosts(rpost.ChannelId)
@@ -484,6 +500,11 @@ func DeletePost(postId string) (*model.Post, *model.AppError) {
go DeletePostFiles(post) go DeletePostFiles(post)
go DeleteFlaggedPosts(post.Id) go DeleteFlaggedPosts(post.Id)
esInterface := einterfaces.GetElasticSearchInterface()
if (esInterface != nil && *utils.Cfg.ElasticSearchSettings.EnableIndexing) {
go esInterface.DeletePost(post.Id)
}
InvalidateCacheForChannelPosts(post.ChannelId) InvalidateCacheForChannelPosts(post.ChannelId)
return post, nil return post, nil
@@ -509,27 +530,84 @@ func DeletePostFiles(post *model.Post) {
func SearchPostsInTeam(terms string, userId string, teamId string, isOrSearch bool) (*model.PostList, *model.AppError) { func SearchPostsInTeam(terms string, userId string, teamId string, isOrSearch bool) (*model.PostList, *model.AppError) {
paramsList := model.ParseSearchParams(terms) paramsList := model.ParseSearchParams(terms)
channels := []store.StoreChannel{}
for _, params := range paramsList { esInterface := einterfaces.GetElasticSearchInterface()
params.OrTerms = isOrSearch if (esInterface != nil && *utils.Cfg.ElasticSearchSettings.EnableSearching && utils.IsLicensed && *utils.License.Features.ElasticSearch) {
// don't allow users to search for everything finalParamsList := []*model.SearchParams{}
if params.Terms != "*" {
channels = append(channels, Srv.Store.Post().Search(teamId, userId, params)) for _, params := range paramsList {
params.OrTerms = isOrSearch
// Don't allow users to search for "*"
if params.Terms != "*" {
// Convert channel names to channel IDs
for idx, channelName := range params.InChannels {
if channel, err := GetChannelByName(channelName, teamId); err != nil {
l4g.Error(err)
} else {
params.InChannels[idx] = channel.Id
}
}
// Convert usernames to user IDs
for idx, username := range params.FromUsers {
if user, err := GetUserByUsername(username); err != nil {
l4g.Error(err)
} else {
params.FromUsers[idx] = user.Id
}
}
finalParamsList = append(finalParamsList, params)
}
} }
}
posts := model.NewPostList() // We only allow the user to search in channels they are a member of.
for _, channel := range channels { userChannels, err := GetChannelsForUser(teamId, userId)
if result := <-channel; result.Err != nil { if err != nil {
return nil, result.Err l4g.Error(err)
return nil, err
}
postIds, err := einterfaces.GetElasticSearchInterface().SearchPosts(userChannels, finalParamsList)
if err != nil {
return nil, err
}
// Get the posts
postList := model.NewPostList()
if presult := <-Srv.Store.Post().GetPostsByIds(postIds); presult.Err != nil {
return nil, presult.Err
} else { } else {
data := result.Data.(*model.PostList) for _, p := range presult.Data.([]*model.Post) {
posts.Extend(data) postList.AddPost(p)
postList.AddOrder(p.Id)
}
} }
}
return posts, nil return postList, nil
} else {
channels := []store.StoreChannel{}
for _, params := range paramsList {
params.OrTerms = isOrSearch
// don't allow users to search for everything
if params.Terms != "*" {
channels = append(channels, Srv.Store.Post().Search(teamId, userId, params))
}
}
posts := model.NewPostList()
for _, channel := range channels {
if result := <-channel; result.Err != nil {
return nil, result.Err
} else {
data := result.Data.(*model.PostList)
posts.Extend(data)
}
}
return posts, nil
}
} }
func GetFileInfosForPost(postId string, readFromMaster bool) ([]*model.FileInfo, *model.AppError) { func GetFileInfosForPost(postId string, readFromMaster bool) ([]*model.FileInfo, *model.AppError) {

Просмотреть файл

@@ -114,6 +114,12 @@ func runServer(configFileLocation string) {
einterfaces.GetMetricsInterface().StartServer() einterfaces.GetMetricsInterface().StartServer()
} }
if einterfaces.GetElasticSearchInterface() != nil {
if err := einterfaces.GetElasticSearchInterface().Start(); err != nil {
l4g.Error(err.Error())
}
}
// wait for kill signal before attempting to gracefully shutdown // wait for kill signal before attempting to gracefully shutdown
// the running service // the running service
c := make(chan os.Signal) c := make(chan os.Signal)

Просмотреть файл

@@ -46,6 +46,14 @@
"EnableUserStatuses": true, "EnableUserStatuses": true,
"ClusterLogTimeoutMilliseconds": 2000 "ClusterLogTimeoutMilliseconds": 2000
}, },
"ElasticSearchSettings": {
"ConnectionUrl": "http://dockerhost:9200",
"Username": "elastic",
"Password": "changeme",
"EnableIndexing": false,
"EnableSearching": false,
"Sniff": true
},
"TeamSettings": { "TeamSettings": {
"SiteName": "Mattermost", "SiteName": "Mattermost",
"MaxUsersPerTeam": 50, "MaxUsersPerTeam": 50,

23
einterfaces/elasticsearch.go Обычный файл
Просмотреть файл

@@ -0,0 +1,23 @@
// Copyright (c) 2017-present Mattermost, Inc. All Rights Reserved.
// See License.txt for license information.
package einterfaces
import "github.com/mattermost/platform/model"
type ElasticSearchInterface interface {
Start() *model.AppError
IndexPost(post *model.Post, teamId string)
SearchPosts(channels *model.ChannelList, searchParams []*model.SearchParams) ([]string, *model.AppError)
DeletePost(postId string)
}
var theElasticSearchInterface ElasticSearchInterface
func RegisterElasticSearchInterface(newInterface ElasticSearchInterface) {
theElasticSearchInterface = newInterface
}
func GetElasticSearchInterface() ElasticSearchInterface {
return theElasticSearchInterface
}

Просмотреть файл

@@ -3307,6 +3307,38 @@
"id": "ent.compliance.run_started.info", "id": "ent.compliance.run_started.info",
"translation": "Compliance export started for job '{{.JobName}}' at '{{.FilePath}}'" "translation": "Compliance export started for job '{{.JobName}}' at '{{.FilePath}}'"
}, },
{
"id": "ent.elasticsearch.start.connect_failed",
"translation": "Setting up ElasticSearch Client Failed"
},
{
"id": "ent.elasticsearch.start.index_exists_failed",
"translation": "Failed to establish whether ElasticSearch index exists"
},
{
"id": "ent.elasticsearch.start.index_create_failed",
"translation": "Failed to create ElasticSearch index"
},
{
"id": "ent.elasticsearch.start.index_settings_failed",
"translation": "Failed to set ElasticSearch index settings"
},
{
"id": "ent.elasticsearch.start.index_mapping_failed",
"translation": "Failed to setup ElasticSearch index mapping"
},
{
"id": "ent.elasticsearch.search_posts.disabled",
"translation": "ElasticSearch searching is disabled on this server"
},
{
"id": "ent.elasticsearch.search_posts.search_failed",
"translation": "Search failed to complete"
},
{
"id": "ent.elasticsearch.search_posts.unmarshall_post_failed",
"translation": "Failed to decode search results"
},
{ {
"id": "ent.emoji.licence_disable.app_error", "id": "ent.emoji.licence_disable.app_error",
"translation": "Custom emoji restrictions disabled by current license. Please contact your system administrator about upgrading your enterprise license." "translation": "Custom emoji restrictions disabled by current license. Please contact your system administrator about upgrading your enterprise license."
@@ -3859,6 +3891,22 @@
"id": "model.compliance.is_valid.start_end_at.app_error", "id": "model.compliance.is_valid.start_end_at.app_error",
"translation": "To must be greater than From" "translation": "To must be greater than From"
}, },
{
"id": "model.config.is_valid.elastic_search.connection_url.app_error",
"translation": "Elastic Search ConnectionUrl setting must be provided when Elastic Search indexing is enabled."
},
{
"id": "model.config.is_valid.elastic_search.username.app_error",
"translation": "Elastic Search Username setting must be provided when Elastic Search indexing is enabled."
},
{
"id": "model.config.is_valid.elastic_search.password.app_error",
"translation": "Elastic Search Password setting must be provided when Elastic Search indexing is enabled."
},
{
"id": "model.config.is_valid.elastic_search.enable_searching.app_error",
"translation": "Elastic Search IndexingEnabled setting must be set to true when Elastic Search SearchEnabled is set to true."
},
{ {
"id": "model.config.is_valid.cluster_email_batching.app_error", "id": "model.config.is_valid.cluster_email_batching.app_error",
"translation": "Unable to enable email batching when clustering is enabled." "translation": "Unable to enable email batching when clustering is enabled."
@@ -5179,6 +5227,10 @@
"id": "store.sql_post.get.app_error", "id": "store.sql_post.get.app_error",
"translation": "We couldn't get the post" "translation": "We couldn't get the post"
}, },
{
"id": "store.sql_post.get_posts_by_ids.app_error",
"translation": "We couldn't get the posts"
},
{ {
"id": "store.sql_post.get_parents_posts.app_error", "id": "store.sql_post.get_parents_posts.app_error",
"translation": "We couldn't get the parent post for the channel" "translation": "We couldn't get the parent post for the channel"

Просмотреть файл

@@ -401,29 +401,39 @@ type WebrtcSettings struct {
TurnSharedKey *string TurnSharedKey *string
} }
type ElasticSearchSettings struct {
ConnectionUrl *string
Username *string
Password *string
EnableIndexing *bool
EnableSearching *bool
Sniff *bool
}
type Config struct { type Config struct {
ServiceSettings ServiceSettings ServiceSettings ServiceSettings
TeamSettings TeamSettings TeamSettings TeamSettings
SqlSettings SqlSettings SqlSettings SqlSettings
LogSettings LogSettings LogSettings LogSettings
PasswordSettings PasswordSettings PasswordSettings PasswordSettings
FileSettings FileSettings FileSettings FileSettings
EmailSettings EmailSettings EmailSettings EmailSettings
RateLimitSettings RateLimitSettings RateLimitSettings RateLimitSettings
PrivacySettings PrivacySettings PrivacySettings PrivacySettings
SupportSettings SupportSettings SupportSettings SupportSettings
GitLabSettings SSOSettings GitLabSettings SSOSettings
GoogleSettings SSOSettings GoogleSettings SSOSettings
Office365Settings SSOSettings Office365Settings SSOSettings
LdapSettings LdapSettings LdapSettings LdapSettings
ComplianceSettings ComplianceSettings ComplianceSettings ComplianceSettings
LocalizationSettings LocalizationSettings LocalizationSettings LocalizationSettings
SamlSettings SamlSettings SamlSettings SamlSettings
NativeAppSettings NativeAppSettings NativeAppSettings NativeAppSettings
ClusterSettings ClusterSettings ClusterSettings ClusterSettings
MetricsSettings MetricsSettings MetricsSettings MetricsSettings
AnalyticsSettings AnalyticsSettings AnalyticsSettings AnalyticsSettings
WebrtcSettings WebrtcSettings WebrtcSettings WebrtcSettings
ElasticSearchSettings ElasticSearchSettings
} }
func (o *Config) ToJson() string { func (o *Config) ToJson() string {
@@ -1217,6 +1227,36 @@ func (o *Config) SetDefaults() {
*o.ServiceSettings.ClusterLogTimeoutMilliseconds = 2000 *o.ServiceSettings.ClusterLogTimeoutMilliseconds = 2000
} }
if o.ElasticSearchSettings.ConnectionUrl == nil {
o.ElasticSearchSettings.ConnectionUrl = new(string)
*o.ElasticSearchSettings.ConnectionUrl = ""
}
if o.ElasticSearchSettings.Username == nil {
o.ElasticSearchSettings.Username = new(string)
*o.ElasticSearchSettings.Username = ""
}
if o.ElasticSearchSettings.Password == nil {
o.ElasticSearchSettings.Password = new(string)
*o.ElasticSearchSettings.Password = ""
}
if o.ElasticSearchSettings.EnableIndexing == nil {
o.ElasticSearchSettings.EnableIndexing = new(bool)
*o.ElasticSearchSettings.EnableIndexing = false
}
if o.ElasticSearchSettings.EnableSearching == nil {
o.ElasticSearchSettings.EnableSearching = new(bool)
*o.ElasticSearchSettings.EnableSearching = false
}
if o.ElasticSearchSettings.Sniff == nil {
o.ElasticSearchSettings.Sniff = new(bool)
*o.ElasticSearchSettings.Sniff = true
}
o.defaultWebrtcSettings() o.defaultWebrtcSettings()
} }
@@ -1448,6 +1488,16 @@ func (o *Config) IsValid() *AppError {
return NewLocAppError("Config.IsValid", "model.config.is_valid.time_between_user_typing.app_error", nil, "") return NewLocAppError("Config.IsValid", "model.config.is_valid.time_between_user_typing.app_error", nil, "")
} }
if *o.ElasticSearchSettings.EnableIndexing {
if len(*o.ElasticSearchSettings.ConnectionUrl) == 0 {
return NewLocAppError("Config.IsValid", "model.config.is_valid.elastic_search.connection_url.app_error", nil, "")
}
}
if *o.ElasticSearchSettings.EnableSearching && !*o.ElasticSearchSettings.EnableIndexing {
return NewLocAppError("Config.IsValid", "model.config.is_valid.elastic_search.enable_searching.app_error", nil, "")
}
return nil return nil
} }
@@ -1488,6 +1538,10 @@ func (o *Config) Sanitize() {
for i := range o.SqlSettings.DataSourceSearchReplicas { for i := range o.SqlSettings.DataSourceSearchReplicas {
o.SqlSettings.DataSourceSearchReplicas[i] = FAKE_SETTING o.SqlSettings.DataSourceSearchReplicas[i] = FAKE_SETTING
} }
*o.ElasticSearchSettings.ConnectionUrl = FAKE_SETTING
*o.ElasticSearchSettings.Username = FAKE_SETTING
*o.ElasticSearchSettings.Password = FAKE_SETTING
} }
func (o *Config) defaultWebrtcSettings() { func (o *Config) defaultWebrtcSettings() {

Просмотреть файл

@@ -1287,3 +1287,30 @@ func (s SqlPostStore) GetPostsCreatedAt(channelId string, time int64) StoreChann
return storeChannel return storeChannel
} }
func (s SqlPostStore) GetPostsByIds(postIds []string) StoreChannel {
storeChannel := make(StoreChannel, 1)
go func() {
result := StoreResult{}
inClause := `'` + strings.Join(postIds, `', '`) + `'`
query := `SELECT * FROM Posts WHERE Id in (` + inClause + `) and DeleteAt = 0 ORDER BY CreateAt DESC`
var posts []*model.Post
_, err := s.GetReplica().Select(&posts, query, map[string]interface{}{})
if err != nil {
l4g.Error(err)
result.Err = model.NewAppError("SqlPostStore.GetPostsCreatedAt", "store.sql_post.get_posts_by_ids.app_error", nil, "", http.StatusInternalServerError)
} else {
result.Data = posts
}
storeChannel <- result
close(storeChannel)
}()
return storeChannel
}

Просмотреть файл

@@ -1550,3 +1550,45 @@ func TestPostStoreOverwrite(t *testing.T) {
t.Fatal("Failed to set FileIds") t.Fatal("Failed to set FileIds")
} }
} }
func TestPostStoreGetPostsByIds(t *testing.T) {
Setup()
o1 := &model.Post{}
o1.ChannelId = model.NewId()
o1.UserId = model.NewId()
o1.Message = "a" + model.NewId() + "AAAAAAAAAAA"
o1 = (<-store.Post().Save(o1)).Data.(*model.Post)
o2 := &model.Post{}
o2.ChannelId = o1.ChannelId
o2.UserId = model.NewId()
o2.Message = "a" + model.NewId() + "CCCCCCCCC"
o2 = (<-store.Post().Save(o2)).Data.(*model.Post)
o3 := &model.Post{}
o3.ChannelId = o1.ChannelId
o3.UserId = model.NewId()
o3.Message = "a" + model.NewId() + "QQQQQQQQQQ"
o3 = (<-store.Post().Save(o3)).Data.(*model.Post)
ro1 := (<-store.Post().Get(o1.Id)).Data.(*model.PostList).Posts[o1.Id]
ro2 := (<-store.Post().Get(o2.Id)).Data.(*model.PostList).Posts[o2.Id]
ro3 := (<-store.Post().Get(o3.Id)).Data.(*model.PostList).Posts[o3.Id]
postIds := []string{
ro1.Id,
ro2.Id,
ro3.Id,
}
if ro4 := Must(store.Post().GetPostsByIds(postIds)).([]*model.Post); len(ro4) != 3 {
t.Fatalf("Expected 3 posts in results. Got %v", len(ro4))
}
Must(store.Post().Delete(ro1.Id, model.GetMillis()))
if ro5 := Must(store.Post().GetPostsByIds(postIds)).([]*model.Post); len(ro5) != 2 {
t.Fatalf("Expected 2 posts in results. Got %v", len(ro5))
}
}

Просмотреть файл

@@ -165,6 +165,7 @@ type PostStore interface {
InvalidateLastPostTimeCache(channelId string) InvalidateLastPostTimeCache(channelId string)
GetPostsCreatedAt(channelId string, time int64) StoreChannel GetPostsCreatedAt(channelId string, time int64) StoreChannel
Overwrite(post *model.Post) StoreChannel Overwrite(post *model.Post) StoreChannel
GetPostsByIds(postIds []string) StoreChannel
} }
type UserStore interface { type UserStore interface {

Просмотреть файл

@@ -492,6 +492,11 @@ func getClientConfig(c *model.Config) map[string]string {
props["PasswordRequireNumber"] = strconv.FormatBool(*c.PasswordSettings.Number) props["PasswordRequireNumber"] = strconv.FormatBool(*c.PasswordSettings.Number)
props["PasswordRequireSymbol"] = strconv.FormatBool(*c.PasswordSettings.Symbol) props["PasswordRequireSymbol"] = strconv.FormatBool(*c.PasswordSettings.Symbol)
} }
if *License.Features.ElasticSearch {
props["ElasticSearchEnableIndexing"] = strconv.FormatBool(*c.ElasticSearchSettings.EnableIndexing)
props["ElasticSearchEnableSearching"] = strconv.FormatBool(*c.ElasticSearchSettings.EnableSearching)
}
} }
return props return props
@@ -560,6 +565,16 @@ func Desanitize(cfg *model.Config) {
cfg.SqlSettings.AtRestEncryptKey = Cfg.SqlSettings.AtRestEncryptKey cfg.SqlSettings.AtRestEncryptKey = Cfg.SqlSettings.AtRestEncryptKey
} }
if *cfg.ElasticSearchSettings.ConnectionUrl == model.FAKE_SETTING {
*cfg.ElasticSearchSettings.ConnectionUrl = *Cfg.ElasticSearchSettings.ConnectionUrl
}
if *cfg.ElasticSearchSettings.Username == model.FAKE_SETTING {
*cfg.ElasticSearchSettings.Username = *Cfg.ElasticSearchSettings.Username
}
if *cfg.ElasticSearchSettings.Password == model.FAKE_SETTING {
*cfg.ElasticSearchSettings.Password = *Cfg.ElasticSearchSettings.Password
}
for i := range cfg.SqlSettings.DataSourceReplicas { for i := range cfg.SqlSettings.DataSourceReplicas {
cfg.SqlSettings.DataSourceReplicas[i] = Cfg.SqlSettings.DataSourceReplicas[i] cfg.SqlSettings.DataSourceReplicas[i] = Cfg.SqlSettings.DataSourceReplicas[i]
} }