From 4932fa058a282d5fbbd2f3c954a6c031343fe272 Mon Sep 17 00:00:00 2001 From: Abhinav Raut Date: Tue, 22 Oct 2024 04:12:13 +0530 Subject: [PATCH] WIP: Multi tag select rule evaluator for contains operator - rename middlewares - Update config.sample.toml --- cmd/handlers.go | 162 +++++++++--------- cmd/init.go | 10 +- cmd/main.go | 8 +- cmd/middlewares.go | 16 +- config.sample.toml | 35 ++-- .../admin/automation/CreateOrEditRule.vue | 13 ++ .../components/admin/automation/RuleBox.vue | 37 +++- internal/automation/evaluator.go | 33 +++- internal/conversation/message.go | 5 +- internal/inbox/inbox.go | 4 +- 10 files changed, 195 insertions(+), 128 deletions(-) diff --git a/cmd/handlers.go b/cmd/handlers.go index 0a149f53..5acd1f7c 100644 --- a/cmd/handlers.go +++ b/cmd/handlers.go @@ -23,126 +23,126 @@ func initHandlers(g *fastglue.Fastglue, hub *ws.Hub) { g.GET("/health", handleHealthCheck) // Serve media files. - g.GET("/uploads/{uuid}", reqAuth(handleServeMedia)) + g.GET("/uploads/{uuid}", auth(handleServeMedia)) // Settings. g.GET("/api/settings/general", handleGetGeneralSettings) - g.PUT("/api/settings/general", reqAuthAndPerm(handleUpdateGeneralSettings, "settings_general", "write")) - g.GET("/api/settings/notifications/email", reqAuthAndPerm(handleGetEmailNotificationSettings, "settings_notifications", "read")) - g.PUT("/api/settings/notifications/email", reqAuthAndPerm(handleUpdateEmailNotificationSettings, "settings_notifications", "write")) + g.PUT("/api/settings/general", authPerm(handleUpdateGeneralSettings, "settings_general", "write")) + g.GET("/api/settings/notifications/email", authPerm(handleGetEmailNotificationSettings, "settings_notifications", "read")) + g.PUT("/api/settings/notifications/email", authPerm(handleUpdateEmailNotificationSettings, "settings_notifications", "write")) // OpenID SSO. g.GET("/api/oidc", handleGetAllOIDC) - g.GET("/api/oidc/{id}", reqAuthAndPerm(handleGetOIDC, "oidc", "read")) - g.POST("/api/oidc", reqAuthAndPerm(handleCreateOIDC, "oidc", "write")) - g.PUT("/api/oidc/{id}", reqAuthAndPerm(handleUpdateOIDC, "oidc", "write")) - g.DELETE("/api/oidc/{id}", reqAuthAndPerm(handleDeleteOIDC, "oidc", "delete")) + g.GET("/api/oidc/{id}", authPerm(handleGetOIDC, "oidc", "read")) + g.POST("/api/oidc", authPerm(handleCreateOIDC, "oidc", "write")) + g.PUT("/api/oidc/{id}", authPerm(handleUpdateOIDC, "oidc", "write")) + g.DELETE("/api/oidc/{id}", authPerm(handleDeleteOIDC, "oidc", "delete")) // Conversation and message. - g.GET("/api/conversations/all", reqAuthAndPerm(handleGetAllConversations, "conversations", "read_all")) - g.GET("/api/conversations/unassigned", reqAuthAndPerm(handleGetUnassignedConversations, "conversations", "read_unassigned")) - g.GET("/api/conversations/assigned", reqAuthAndPerm(handleGetAssignedConversations, "conversations", "read_assigned")) - g.GET("/api/conversations/{uuid}", reqAuthAndPerm(handleGetConversation, "conversations", "read")) - g.GET("/api/conversations/{uuid}/participants", reqAuthAndPerm(handleGetConversationParticipants, "conversations", "read")) - g.PUT("/api/conversations/{uuid}/assignee/user", reqAuthAndPerm(handleUpdateConversationUserAssignee, "conversations", "update_user_assignee")) - g.PUT("/api/conversations/{uuid}/assignee/team", reqAuthAndPerm(handleUpdateTeamAssignee, "conversations", "update_team_assignee")) - g.PUT("/api/conversations/{uuid}/priority", reqAuthAndPerm(handleUpdateConversationPriority, "conversations", "update_priority")) - g.PUT("/api/conversations/{uuid}/status", reqAuthAndPerm(handleUpdateConversationStatus, "conversations", "update_status")) - g.PUT("/api/conversations/{uuid}/last-seen", reqAuthAndPerm(handleUpdateConversationAssigneeLastSeen, "conversations", "read")) - g.POST("/api/conversations/{uuid}/tags", reqAuthAndPerm(handleAddConversationTags, "conversations", "update_tags")) - g.POST("/api/conversations/{cuuid}/messages", reqAuthAndPerm(handleSendMessage, "messages", "write")) - g.GET("/api/conversations/{uuid}/messages", reqAuthAndPerm(handleGetMessages, "messages", "read")) - g.PUT("/api/conversations/{cuuid}/messages/{uuid}/retry", reqAuthAndPerm(handleRetryMessage, "messages", "write")) - g.GET("/api/conversations/{cuuid}/messages/{uuid}", reqAuthAndPerm(handleGetMessage, "messages", "read")) + g.GET("/api/conversations/all", authPerm(handleGetAllConversations, "conversations", "read_all")) + g.GET("/api/conversations/unassigned", authPerm(handleGetUnassignedConversations, "conversations", "read_unassigned")) + g.GET("/api/conversations/assigned", authPerm(handleGetAssignedConversations, "conversations", "read_assigned")) + g.GET("/api/conversations/{uuid}", authPerm(handleGetConversation, "conversations", "read")) + g.GET("/api/conversations/{uuid}/participants", authPerm(handleGetConversationParticipants, "conversations", "read")) + g.PUT("/api/conversations/{uuid}/assignee/user", authPerm(handleUpdateConversationUserAssignee, "conversations", "update_user_assignee")) + g.PUT("/api/conversations/{uuid}/assignee/team", authPerm(handleUpdateTeamAssignee, "conversations", "update_team_assignee")) + g.PUT("/api/conversations/{uuid}/priority", authPerm(handleUpdateConversationPriority, "conversations", "update_priority")) + g.PUT("/api/conversations/{uuid}/status", authPerm(handleUpdateConversationStatus, "conversations", "update_status")) + g.PUT("/api/conversations/{uuid}/last-seen", authPerm(handleUpdateConversationAssigneeLastSeen, "conversations", "read")) + g.POST("/api/conversations/{uuid}/tags", authPerm(handleAddConversationTags, "conversations", "update_tags")) + g.POST("/api/conversations/{cuuid}/messages", authPerm(handleSendMessage, "messages", "write")) + g.GET("/api/conversations/{uuid}/messages", authPerm(handleGetMessages, "messages", "read")) + g.PUT("/api/conversations/{cuuid}/messages/{uuid}/retry", authPerm(handleRetryMessage, "messages", "write")) + g.GET("/api/conversations/{cuuid}/messages/{uuid}", authPerm(handleGetMessage, "messages", "read")) // Status and priority. - g.GET("/api/statuses", reqAuth(handleGetStatuses)) - g.POST("/api/statuses", reqAuthAndPerm(handleCreateStatus, "status", "write")) - g.PUT("/api/statuses/{id}", reqAuthAndPerm(handleUpdateStatus, "status", "write")) - g.DELETE("/api/statuses/{id}", reqAuthAndPerm(handleDeleteStatus, "status", "delete")) - g.GET("/api/priorities", reqAuth(handleGetPriorities)) + g.GET("/api/statuses", auth(handleGetStatuses)) + g.POST("/api/statuses", authPerm(handleCreateStatus, "status", "write")) + g.PUT("/api/statuses/{id}", authPerm(handleUpdateStatus, "status", "write")) + g.DELETE("/api/statuses/{id}", authPerm(handleDeleteStatus, "status", "delete")) + g.GET("/api/priorities", auth(handleGetPriorities)) // Tag. - g.GET("/api/tags", reqAuth(handleGetTags)) - g.POST("/api/tags", reqAuthAndPerm(handleCreateTag, "tags", "write")) - g.PUT("/api/tags/{id}", reqAuthAndPerm(handleUpdateTag, "tags", "write")) - g.DELETE("/api/tags/{id}", reqAuthAndPerm(handleDeleteTag, "tags", "delete")) + g.GET("/api/tags", auth(handleGetTags)) + g.POST("/api/tags", authPerm(handleCreateTag, "tags", "write")) + g.PUT("/api/tags/{id}", authPerm(handleUpdateTag, "tags", "write")) + g.DELETE("/api/tags/{id}", authPerm(handleDeleteTag, "tags", "delete")) // Media. - g.POST("/api/media", reqAuth(handleMediaUpload)) + g.POST("/api/media", auth(handleMediaUpload)) // Canned response. - g.GET("/api/canned-responses", reqAuth(handleGetCannedResponses)) - g.POST("/api/canned-responses", reqAuthAndPerm(handleCreateCannedResponse, "canned_responses", "write")) - g.PUT("/api/canned-responses/{id}", reqAuthAndPerm(handleUpdateCannedResponse, "canned_responses", "write")) - g.DELETE("/api/canned-responses/{id}", reqAuthAndPerm(handleDeleteCannedResponse, "canned_responses", "delete")) + g.GET("/api/canned-responses", auth(handleGetCannedResponses)) + g.POST("/api/canned-responses", authPerm(handleCreateCannedResponse, "canned_responses", "write")) + g.PUT("/api/canned-responses/{id}", authPerm(handleUpdateCannedResponse, "canned_responses", "write")) + g.DELETE("/api/canned-responses/{id}", authPerm(handleDeleteCannedResponse, "canned_responses", "delete")) // User. - g.GET("/api/users/me", reqAuth(handleGetCurrentUser)) - g.PUT("/api/users/me", reqAuth(handleUpdateCurrentUser)) - g.DELETE("/api/users/me/avatar", reqAuth(handleDeleteAvatar)) - g.GET("/api/users/compact", reqAuth(handleGetUsersCompact)) - g.GET("/api/users", reqAuthAndPerm(handleGetUsers, "users", "read")) - g.GET("/api/users/{id}", reqAuthAndPerm(handleGetUser, "users", "read")) - g.POST("/api/users", reqAuthAndPerm(handleCreateUser, "users", "write")) - g.PUT("/api/users/{id}", reqAuthAndPerm(handleUpdateUser, "users", "write")) + g.GET("/api/users/me", auth(handleGetCurrentUser)) + g.PUT("/api/users/me", auth(handleUpdateCurrentUser)) + g.DELETE("/api/users/me/avatar", auth(handleDeleteAvatar)) + g.GET("/api/users/compact", auth(handleGetUsersCompact)) + g.GET("/api/users", authPerm(handleGetUsers, "users", "read")) + g.GET("/api/users/{id}", authPerm(handleGetUser, "users", "read")) + g.POST("/api/users", authPerm(handleCreateUser, "users", "write")) + g.PUT("/api/users/{id}", authPerm(handleUpdateUser, "users", "write")) // Team. - g.GET("/api/teams/compact", reqAuth(handleGetTeamsCompact)) - g.GET("/api/teams", reqAuthAndPerm(handleGetTeams, "teams", "read")) - g.GET("/api/teams/{id}", reqAuthAndPerm(handleGetTeam, "teams", "read")) - g.PUT("/api/teams/{id}", reqAuthAndPerm(handleUpdateTeam, "teams", "write")) - g.POST("/api/teams", reqAuthAndPerm(handleCreateTeam, "teams", "write")) + g.GET("/api/teams/compact", auth(handleGetTeamsCompact)) + g.GET("/api/teams", authPerm(handleGetTeams, "teams", "read")) + g.GET("/api/teams/{id}", authPerm(handleGetTeam, "teams", "read")) + g.PUT("/api/teams/{id}", authPerm(handleUpdateTeam, "teams", "write")) + g.POST("/api/teams", authPerm(handleCreateTeam, "teams", "write")) // i18n. g.GET("/api/lang/{lang}", handleGetI18nLang) // Automation. - g.GET("/api/automation/rules", reqAuthAndPerm(handleGetAutomationRules, "automations", "read")) - g.GET("/api/automation/rules/{id}", reqAuthAndPerm(handleGetAutomationRule, "automations", "read")) - g.POST("/api/automation/rules", reqAuthAndPerm(handleCreateAutomationRule, "automations", "write")) - g.PUT("/api/automation/rules/{id}/toggle", reqAuthAndPerm(handleToggleAutomationRule, "automations", "write")) - g.PUT("/api/automation/rules/{id}", reqAuthAndPerm(handleUpdateAutomationRule, "automations", "write")) - g.DELETE("/api/automation/rules/{id}", reqAuthAndPerm(handleDeleteAutomationRule, "automations", "delete")) + g.GET("/api/automation/rules", authPerm(handleGetAutomationRules, "automations", "read")) + g.GET("/api/automation/rules/{id}", authPerm(handleGetAutomationRule, "automations", "read")) + g.POST("/api/automation/rules", authPerm(handleCreateAutomationRule, "automations", "write")) + g.PUT("/api/automation/rules/{id}/toggle", authPerm(handleToggleAutomationRule, "automations", "write")) + g.PUT("/api/automation/rules/{id}", authPerm(handleUpdateAutomationRule, "automations", "write")) + g.DELETE("/api/automation/rules/{id}", authPerm(handleDeleteAutomationRule, "automations", "delete")) // Inbox. - g.GET("/api/inboxes", reqAuthAndPerm(handleGetInboxes, "inboxes", "read")) - g.GET("/api/inboxes/{id}", reqAuthAndPerm(handleGetInbox, "inboxes", "read")) - g.POST("/api/inboxes", reqAuthAndPerm(handleCreateInbox, "inboxes", "write")) - g.PUT("/api/inboxes/{id}/toggle", reqAuthAndPerm(handleToggleInbox, "inboxes", "write")) - g.PUT("/api/inboxes/{id}", reqAuthAndPerm(handleUpdateInbox, "inboxes", "write")) - g.DELETE("/api/inboxes/{id}", reqAuthAndPerm(handleDeleteInbox, "inboxes", "delete")) + g.GET("/api/inboxes", authPerm(handleGetInboxes, "inboxes", "read")) + g.GET("/api/inboxes/{id}", authPerm(handleGetInbox, "inboxes", "read")) + g.POST("/api/inboxes", authPerm(handleCreateInbox, "inboxes", "write")) + g.PUT("/api/inboxes/{id}/toggle", authPerm(handleToggleInbox, "inboxes", "write")) + g.PUT("/api/inboxes/{id}", authPerm(handleUpdateInbox, "inboxes", "write")) + g.DELETE("/api/inboxes/{id}", authPerm(handleDeleteInbox, "inboxes", "delete")) // Role. - g.GET("/api/roles", reqAuthAndPerm(handleGetRoles, "roles", "read")) - g.GET("/api/roles/{id}", reqAuthAndPerm(handleGetRole, "roles", "read")) - g.POST("/api/roles", reqAuthAndPerm(handleCreateRole, "roles", "write")) - g.PUT("/api/roles/{id}", reqAuthAndPerm(handleUpdateRole, "roles", "write")) - g.DELETE("/api/roles/{id}", reqAuthAndPerm(handleDeleteRole, "roles", "delete")) + g.GET("/api/roles", authPerm(handleGetRoles, "roles", "read")) + g.GET("/api/roles/{id}", authPerm(handleGetRole, "roles", "read")) + g.POST("/api/roles", authPerm(handleCreateRole, "roles", "write")) + g.PUT("/api/roles/{id}", authPerm(handleUpdateRole, "roles", "write")) + g.DELETE("/api/roles/{id}", authPerm(handleDeleteRole, "roles", "delete")) // Dashboard. - g.GET("/api/dashboard/global/counts", reqAuthAndPerm(handleDashboardCounts, "dashboard_global", "read")) - g.GET("/api/dashboard/global/charts", reqAuthAndPerm(handleDashboardCharts, "dashboard_global", "read")) + g.GET("/api/dashboard/global/counts", authPerm(handleDashboardCounts, "dashboard_global", "read")) + g.GET("/api/dashboard/global/charts", authPerm(handleDashboardCharts, "dashboard_global", "read")) // Template. - g.GET("/api/templates", reqAuthAndPerm(handleGetTemplates, "templates", "read")) - g.GET("/api/templates/{id}", reqAuthAndPerm(handleGetTemplate, "templates", "read")) - g.POST("/api/templates", reqAuthAndPerm(handleCreateTemplate, "templates", "write")) - g.PUT("/api/templates/{id}", reqAuthAndPerm(handleUpdateTemplate, "templates", "write")) - g.DELETE("/api/templates/{id}", reqAuthAndPerm(handleDeleteTemplate, "templates", "delete")) + g.GET("/api/templates", authPerm(handleGetTemplates, "templates", "read")) + g.GET("/api/templates/{id}", authPerm(handleGetTemplate, "templates", "read")) + g.POST("/api/templates", authPerm(handleCreateTemplate, "templates", "write")) + g.PUT("/api/templates/{id}", authPerm(handleUpdateTemplate, "templates", "write")) + g.DELETE("/api/templates/{id}", authPerm(handleDeleteTemplate, "templates", "delete")) // WebSocket. - g.GET("/api/ws", reqAuth(func(r *fastglue.Request) error { + g.GET("/api/ws", auth(func(r *fastglue.Request) error { return handleWS(r, hub) })) // Frontend pages. - g.GET("/", notAuthenticatedPage(serveIndexPage)) - g.GET("/dashboard", authenticatedPage(serveIndexPage)) - g.GET("/conversations", authenticatedPage(serveIndexPage)) - g.GET("/conversations/{all:*}", authenticatedPage(serveIndexPage)) - g.GET("/account/profile", authenticatedPage(serveIndexPage)) - g.GET("/admin/{all:*}", authenticatedPage(serveIndexPage)) + g.GET("/", notAuthPage(serveIndexPage)) + g.GET("/dashboard", authPage(serveIndexPage)) + g.GET("/conversations", authPage(serveIndexPage)) + g.GET("/conversations/{all:*}", authPage(serveIndexPage)) + g.GET("/account/profile", authPage(serveIndexPage)) + g.GET("/admin/{all:*}", authPage(serveIndexPage)) g.GET("/assets/{all:*}", serveStaticFiles) g.GET("/images/{all:*}", serveStaticFiles) } diff --git a/cmd/init.go b/cmd/init.go index 469c60aa..270f292e 100644 --- a/cmd/init.go +++ b/cmd/init.go @@ -8,7 +8,7 @@ import ( "os" "path/filepath" - "github.com/abhinavxd/artemis/internal/auth" + auth_ "github.com/abhinavxd/artemis/internal/auth" "github.com/abhinavxd/artemis/internal/authz" "github.com/abhinavxd/artemis/internal/autoassigner" "github.com/abhinavxd/artemis/internal/automation" @@ -431,7 +431,7 @@ func initAuthz() *authz.Enforcer { } // initAuth initializes authentication manager. -func initAuth(o *oidc.Manager, rd *redis.Client) *auth.Auth { +func initAuth(o *oidc.Manager, rd *redis.Client) *auth_.Auth { var lo = initLogger("auth") oidc, err := o.GetAll() @@ -439,12 +439,12 @@ func initAuth(o *oidc.Manager, rd *redis.Client) *auth.Auth { log.Fatalf("error initializing auth: %v", err) } - var providers = make([]auth.Provider, 0, len(oidc)) + var providers = make([]auth_.Provider, 0, len(oidc)) for _, o := range oidc { if o.Disabled { continue } - providers = append(providers, auth.Provider{ + providers = append(providers, auth_.Provider{ ID: o.ID, Provider: o.Provider, ProviderURL: o.ProviderURL, @@ -454,7 +454,7 @@ func initAuth(o *oidc.Manager, rd *redis.Client) *auth.Auth { }) } - auth, err := auth.New(auth.Config{ + auth, err := auth_.New(auth_.Config{ Providers: providers, }, rd, lo) if err != nil { diff --git a/cmd/main.go b/cmd/main.go index 41fdc4b2..fa91b27c 100644 --- a/cmd/main.go +++ b/cmd/main.go @@ -7,7 +7,7 @@ import ( "os/signal" "syscall" - "github.com/abhinavxd/artemis/internal/auth" + auth_ "github.com/abhinavxd/artemis/internal/auth" "github.com/abhinavxd/artemis/internal/authz" "github.com/abhinavxd/artemis/internal/automation" @@ -39,7 +39,7 @@ var ko = koanf.New(".") type App struct { constant constants fs stuffbin.FileSystem - auth *auth.Auth + auth *auth_.Auth authz *authz.Enforcer i18n *i18n.I18n lo *logf.Logger @@ -130,7 +130,7 @@ func main() { automation.SetConversationStore(conversation) // Start receivers for each inbox. - go inbox.Receive(ctx) + go inbox.Start(ctx) // Start evaluating automation rules. go automation.Run(ctx, automationWrk) @@ -139,7 +139,7 @@ func main() { go autoassigner.Run(ctx) // Start listening and dispatching messages. - go conversation.ListenAndDispatchMessages(ctx, messageDispatchWrk, messageDispatchScanInterval) + go conversation.Run(ctx, messageDispatchWrk, messageDispatchScanInterval) // Start notification service. go notifier.Run(ctx) diff --git a/cmd/middlewares.go b/cmd/middlewares.go index 044ea44b..1743583d 100644 --- a/cmd/middlewares.go +++ b/cmd/middlewares.go @@ -8,8 +8,8 @@ import ( "github.com/zerodha/fastglue" ) -// reqAuth makes sure the user is logged in. -func reqAuth(handler fastglue.FastRequestHandler) fastglue.FastRequestHandler { +// auth makes sure the user is logged in. +func auth(handler fastglue.FastRequestHandler) fastglue.FastRequestHandler { return func(r *fastglue.Request) error { var ( app = r.Context.(*App) @@ -33,8 +33,8 @@ func reqAuth(handler fastglue.FastRequestHandler) fastglue.FastRequestHandler { } } -// reqAuthAndPerm does session validation, CSRF, and permission enforcement. -func reqAuthAndPerm(handler fastglue.FastRequestHandler, object, action string) fastglue.FastRequestHandler { +// authPerm does session validation, CSRF, and permission enforcement. +func authPerm(handler fastglue.FastRequestHandler, object, action string) fastglue.FastRequestHandler { return func(r *fastglue.Request) error { var ( app = r.Context.(*App) @@ -77,8 +77,8 @@ func reqAuthAndPerm(handler fastglue.FastRequestHandler, object, action string) } } -// authenticatedPage ensures the user is logged in; otherwise, redirects to the login page. -func authenticatedPage(handler fastglue.FastRequestHandler) fastglue.FastRequestHandler { +// authPage ensures the user is logged in; otherwise, redirects to the login page. +func authPage(handler fastglue.FastRequestHandler) fastglue.FastRequestHandler { return func(r *fastglue.Request) error { app := r.Context.(*App) @@ -102,8 +102,8 @@ func authenticatedPage(handler fastglue.FastRequestHandler) fastglue.FastRequest } } -// notAuthenticatedPage allows access only if the user is not authenticated; otherwise, redirects to the dashboard. -func notAuthenticatedPage(handler fastglue.FastRequestHandler) fastglue.FastRequestHandler { +// notAuthPage allows access only if the user is not authenticated; otherwise, redirects to the dashboard. +func notAuthPage(handler fastglue.FastRequestHandler) fastglue.FastRequestHandler { return func(r *fastglue.Request) error { app := r.Context.(*App) diff --git a/config.sample.toml b/config.sample.toml index e07559be..8ec4cc36 100644 --- a/config.sample.toml +++ b/config.sample.toml @@ -5,7 +5,7 @@ env = "dev" # HTTP server. [app.server] -name = "artemis" +name = "" address = "0.0.0.0:9009" socket = "" read_timeout = "5s" @@ -13,46 +13,53 @@ write_timeout = "5s" max_body_size = 10000000 keepalive_timeout = "10s" -# Upload. +# File upload. [upload] -provider = "s3" +provider = "fs" -[upload.localfs] -upload_path = "" +# Filesytem provider. +[upload.fs] +upload_path = '/home/ubuntu/uploads' +# S3 provider. [upload.s3] url = "" public_url = "" access_key = "" secret_key = "" region = "ap-south-1" -bucket = "artemis" +bucket = "bucket" bucket_path = "" bucket_type = "private" expiry = "15m" +# Postgres. [db] host = "127.0.0.1" port = 5432 user = "postgres" password = "postgres" -database = "artemis" +database = "database" ssl_mode = "disable" -max_open = 100 -max_idle = 30 +max_open = 10 +max_idle = 10 max_lifetime = "10s" +# Redis. [redis] address = "127.0.0.1:6379" password = "" db = 0 [message] -dispatch_concurrency = 10 -dispatch_read_interval = "50ms" -incoming_queue_size = 10000 -outgoing_queue_size = 10000 +dispatch_workers = 10 +dispatch_scan_interval = "50ms" +incoming_queue_size = 5000 +outgoing_queue_size = 5000 [notification] concurrency = 2 -queue_size = 100 \ No newline at end of file +queue_size = 100 + +[automation] +worker_count = 10 diff --git a/frontend/src/components/admin/automation/CreateOrEditRule.vue b/frontend/src/components/admin/automation/CreateOrEditRule.vue index 9fdb6e2e..8021d7b3 100644 --- a/frontend/src/components/admin/automation/CreateOrEditRule.vue +++ b/frontend/src/components/admin/automation/CreateOrEditRule.vue @@ -305,7 +305,20 @@ onMounted(async () => { } } firstRuleGroup.value = getFirstGroup() + // Convert multi tag select values separated by commas to an array + firstRuleGroup.value.rules.forEach(rule => { + if (!Array.isArray(rule.value)) + if (["contains", "not contains"].includes(rule.operator)) { + rule.value = rule.value ? rule.value.split(',') : [] + } + }) secondRuleGroup.value = getSecondGroup() + secondRuleGroup.value?.rules?.forEach(rule => { + if (!Array.isArray(rule.value)) + if (["contains", "not contains"].includes(rule.operator)) { + rule.value = rule.value ? rule.value.split(',') : [] + } + }) groupOperator.value = getGroupOperator() }) diff --git a/frontend/src/components/admin/automation/RuleBox.vue b/frontend/src/components/admin/automation/RuleBox.vue index 2112414e..888f6cca 100644 --- a/frontend/src/components/admin/automation/RuleBox.vue +++ b/frontend/src/components/admin/automation/RuleBox.vue @@ -57,7 +57,7 @@
- @@ -74,6 +74,15 @@ + + + + + + + + +
@@ -104,6 +113,7 @@ import { SelectTrigger, SelectValue } from '@/components/ui/select' +import { TagsInput, TagsInputInput, TagsInputItem, TagsInputItemDelete, TagsInputItemText } from '@/components/ui/tags-input' import { CircleX } from 'lucide-vue-next' import { Label } from '@/components/ui/label' import { Input } from '@/components/ui/input' @@ -159,17 +169,31 @@ const handleFieldChange = (value, ruleIndex) => { } const handleOperatorChange = (value, ruleIndex) => { - // Clear value on operator change - ruleGroup.value.rules[ruleIndex].value = '' + // Set initial value based on operator. + if (["contains", "not contains"].includes(value)) + ruleGroup.value.rules[ruleIndex].value = [] + else + ruleGroup.value.rules[ruleIndex].value = '' ruleGroup.value.rules[ruleIndex].operator = value emitUpdate() } const handleValueChange = (value, ruleIndex) => { - ruleGroup.value.rules[ruleIndex].value = value + console.log("value ", value) + const operator = ruleGroup.value.rules[ruleIndex].operator + // For 'contains' and 'not contains', join array into a single string + if (["contains", "not contains"].includes(operator)) { + ruleGroup.value.rules[ruleIndex].value = Array.isArray(value) ? value.join(',') : value + } else { + ruleGroup.value.rules[ruleIndex].value = String(value) + } emitUpdate() } +const fieldValueAsArray = (value) => { + return Array.isArray(value) ? value : (value ? value.split(',') : []) +} + const handleCaseSensitiveCheck = (value, ruleIndex) => { ruleGroup.value.rules[ruleIndex].case_sensitive_match = value emitUpdate() @@ -221,10 +245,13 @@ const getFieldOptions = (field) => { const inputType = (index) => { const field = ruleGroup.value.rules[index]?.field const operator = ruleGroup.value.rules[index]?.operator + if (["contains", "not contains"].includes(operator)) { + return "tag" + } if (["status", "priority"].includes(field)) { return "select" } - if (["equals", "not equals", "contains", "not contains"].includes(operator)) { + if (["equals", "not equals"].includes(operator)) { return "text" } return "" diff --git a/internal/automation/evaluator.go b/internal/automation/evaluator.go index 36a7a83f..55f17414 100644 --- a/internal/automation/evaluator.go +++ b/internal/automation/evaluator.go @@ -92,8 +92,9 @@ func (e *Engine) evaluateGroup(rules []models.RuleDetail, operator string, conve // evaluateRule evaluates a single rule against a given conversation. func (e *Engine) evaluateRule(rule models.RuleDetail, conversation cmodels.Conversation) bool { var ( - valueToCompare string - conditionMet bool + valueToCompare string + valuesToCompare []string + conditionMet bool ) // Extract the value from the conversation based on the rule's field @@ -124,7 +125,16 @@ func (e *Engine) evaluateRule(rule models.RuleDetail, conversation cmodels.Conve rule.Value = strings.ToLower(rule.Value) } - e.lo.Debug("evaluating rule", "rule_field", rule.Field, "rule_operator", rule.Operator, "rule_value", rule.Value, "compared_with", valueToCompare, "coversation_uuid", conversation.UUID) + // RuleContains and NotContains have all values comma-separated. + if rule.Operator == models.RuleContains || rule.Operator == models.RuleNotContains { + valuesToCompare = strings.Split(rule.Value, ",") + // Trim whitespace from each value + for i := range valuesToCompare { + valuesToCompare[i] = strings.TrimSpace(valuesToCompare[i]) + } + } + + e.lo.Debug("evaluating rule", "rule_field", rule.Field, "rule_operator", rule.Operator, "rule_value", rule.Value, "values_to_compare", valuesToCompare, "value_to_compare", valueToCompare, "conversation_uuid", conversation.UUID) // Compare with set operator. switch rule.Operator { @@ -133,9 +143,20 @@ func (e *Engine) evaluateRule(rule models.RuleDetail, conversation cmodels.Conve case models.RuleNotEqual: conditionMet = valueToCompare != rule.Value case models.RuleContains: - conditionMet = strings.Contains(valueToCompare, rule.Value) + for _, val := range valuesToCompare { + if strings.Contains(valueToCompare, val) { + conditionMet = true + break + } + } case models.RuleNotContains: - conditionMet = !strings.Contains(valueToCompare, rule.Value) + conditionMet = true + for _, val := range valuesToCompare { + if strings.Contains(valueToCompare, val) { + conditionMet = false + break + } + } case models.RuleSet: conditionMet = len(valueToCompare) > 0 case models.RuleNotSet: @@ -144,7 +165,7 @@ func (e *Engine) evaluateRule(rule models.RuleDetail, conversation cmodels.Conve e.lo.Error("unrecognized rule logical operator", "operator", rule.Operator) return false } - e.lo.Debug("rule conditions met status", "met", conditionMet, "coversation_uuid", conversation.UUID) + e.lo.Debug("rule conditions met status", "met", conditionMet, "conversation_uuid", conversation.UUID) return conditionMet } diff --git a/internal/conversation/message.go b/internal/conversation/message.go index e741bef5..ed8f4373 100644 --- a/internal/conversation/message.go +++ b/internal/conversation/message.go @@ -48,10 +48,9 @@ const ( maxMessagesPerPage = 30 ) -// ListenAndDispatchMessages starts a pool of worker goroutines to handle -// message dispatching via inbox's channel and processes incoming messages. It scans for +// Run starts a pool of worker goroutines to handle message dispatching via inbox's channel and processes incoming messages. It scans for // pending outgoing messages at the specified read interval and pushes them to the outgoing queue. -func (m *Manager) ListenAndDispatchMessages(ctx context.Context, dispatchConcurrency int, scanInterval time.Duration) { +func (m *Manager) Run(ctx context.Context, dispatchConcurrency int, scanInterval time.Duration) { dbScanner := time.NewTicker(scanInterval) defer dbScanner.Stop() diff --git a/internal/inbox/inbox.go b/internal/inbox/inbox.go index d934e751..8f36e7ff 100644 --- a/internal/inbox/inbox.go +++ b/internal/inbox/inbox.go @@ -182,8 +182,8 @@ func (m *Manager) Delete(id int) error { return nil } -// Receive starts the receiver for each inbox. -func (m *Manager) Receive(ctx context.Context) error { +// Start starts the receiver for each inbox. +func (m *Manager) Start(ctx context.Context) error { for _, inb := range m.inboxes { m.wg.Add(1) go func(inbox Inbox) {