mirror of
https://github.com/memohai/Memoh.git
synced 2026-04-27 07:16:19 +09:00
refactor: memory provider (#140)
* refactor: memory provider * fix: migrations * feat: divide collection from different built-in memory * feat: add `MEMORY.md` and `PROFILES.md` * use .env for docker compose. fix #142 (#143) * feat(web): add brand icons for search providers (#144) Add custom FontAwesome icon definitions for all 9 search providers: - Yandex: uses existing faYandex from FA free brands - Tavily, Jina, Exa, Bocha, Serper: custom icons from brand SVGs - DuckDuckGo, SearXNG, Sogou: custom icons from Simple Icons Icons are registered with a custom 'fac' prefix and rendered as monochrome (currentColor) via FontAwesome's standard rendering. * fix: resolve multiple UI bugs (#147) * feat: add email service with multi-adapter support (#146) * feat: add email service with multi-adapter support Implement a full-stack email service with global provider management, per-bot bindings with granular read/write permissions, outbox audit storage, and MCP tool integration for direct mailbox access. Backend: - Email providers: CRUD with dynamic config schema (generic SMTP/IMAP, Mailgun) - Generic adapter: go-mail (SMTP) + go-imap/v2 (IMAP IDLE real-time push via UnilateralDataHandler + UID-based tracking + periodic check fallback) - Mailgun adapter: mailgun-go/v5 with dual inbound mode (webhook + poll) - Bot email bindings: per-bot provider binding with independent r/w permissions - Outbox: outbound email audit log with status tracking - Trigger: inbound emails push notification to bot_inbox (from/subject only, LLM reads full content on demand via MCP tools) - MailboxReader interface: on-demand IMAP queries for listing/reading emails - MCP tools: email_accounts, email_send, email_list (paginated mailbox), email_read (by UID) — all with multi-binding and provider_id selection - Webhook: /email/mailgun/webhook/:config_id (JWT-skipped, signature-verified) - DB migration: 0019_add_email (email_providers, bot_email_bindings, email_outbox) Frontend: - Email Providers page: /email-providers with MasterDetailSidebarLayout - Dynamic config form rendered from ordered provider meta schema with i18n keys - Bot detail: Email tab with bindings management + outbox audit table - Sidebar navigation entry - Full i18n support (en + zh) - Auto-generated SDK from Swagger Closes #17 * feat(email): trigger bot conversation immediately on inbound email Instead of only storing an inbox item and waiting for the next chat, the email trigger now proactively invokes the conversation resolver so the bot processes new emails right away — aligned with the schedule/heartbeat trigger pattern. * fix: lint --------- Co-authored-by: Acbox <acbox0328@gmail.com> * chore: update AGENTS.md * feat: files preview * feat(web): improve MCP details page * refactor(skills): import skill with pure markdown string * merge main into refactor/memory * fix: migration * refactor: temp delete qdrant and bm25 index * fix: clean merge code * fix: update memory handler --------- Co-authored-by: Leohearts <leohearts@leohearts.com> Co-authored-by: Menci <mencici@msn.com> Co-authored-by: Quincy <69751197+dqygit@users.noreply.github.com> Co-authored-by: BBQ <35603386+HoneyBBQ@users.noreply.github.com> Co-authored-by: Ran <16112591+chen-ran@users.noreply.github.com>
This commit is contained in:
+46
-104
@@ -37,7 +37,6 @@ import (
|
||||
"github.com/memohai/memoh/internal/conversation/flow"
|
||||
"github.com/memohai/memoh/internal/db"
|
||||
dbsqlc "github.com/memohai/memoh/internal/db/sqlc"
|
||||
"github.com/memohai/memoh/internal/embeddings"
|
||||
"github.com/memohai/memoh/internal/handlers"
|
||||
"github.com/memohai/memoh/internal/healthcheck"
|
||||
channelchecker "github.com/memohai/memoh/internal/healthcheck/checkers/channel"
|
||||
@@ -56,7 +55,7 @@ import (
|
||||
mcpweb "github.com/memohai/memoh/internal/mcp/providers/web"
|
||||
mcpfederation "github.com/memohai/memoh/internal/mcp/sources/federation"
|
||||
"github.com/memohai/memoh/internal/media"
|
||||
"github.com/memohai/memoh/internal/memory"
|
||||
memprovider "github.com/memohai/memoh/internal/memory/provider"
|
||||
"github.com/memohai/memoh/internal/message"
|
||||
"github.com/memohai/memoh/internal/message/event"
|
||||
"github.com/memohai/memoh/internal/models"
|
||||
@@ -145,12 +144,8 @@ func runServe() {
|
||||
|
||||
// memory pipeline
|
||||
provideMemoryLLM,
|
||||
provideEmbeddingsResolver,
|
||||
provideEmbeddingSetup,
|
||||
provideTextEmbedderForMemory,
|
||||
provideQdrantStore,
|
||||
memory.NewBM25Indexer,
|
||||
provideMemoryService,
|
||||
memprovider.NewService,
|
||||
provideMemoryProviderRegistry,
|
||||
|
||||
// domain services (auto-wired)
|
||||
models.NewService,
|
||||
@@ -205,7 +200,6 @@ func runServe() {
|
||||
provideServerHandler(handlers.NewPingHandler),
|
||||
provideServerHandler(provideAuthHandler),
|
||||
provideServerHandler(provideMemoryHandler),
|
||||
provideServerHandler(handlers.NewEmbeddingsHandler),
|
||||
provideServerHandler(provideMessageHandler),
|
||||
provideServerHandler(handlers.NewSwaggerHandler),
|
||||
provideServerHandler(handlers.NewProvidersHandler),
|
||||
@@ -220,6 +214,7 @@ func runServe() {
|
||||
provideServerHandler(handlers.NewChannelHandler),
|
||||
provideServerHandler(feishu.NewWebhookServerHandler),
|
||||
provideServerHandler(provideUsersHandler),
|
||||
provideServerHandler(handlers.NewMemoryProvidersHandler),
|
||||
provideServerHandler(handlers.NewEmailProvidersHandler),
|
||||
provideServerHandler(handlers.NewEmailBindingsHandler),
|
||||
provideServerHandler(handlers.NewEmailOutboxHandler),
|
||||
@@ -233,7 +228,7 @@ func runServe() {
|
||||
provideServer,
|
||||
),
|
||||
fx.Invoke(
|
||||
startMemoryWarmup,
|
||||
startMemoryProviderBootstrap,
|
||||
startScheduleService,
|
||||
startHeartbeatService,
|
||||
startChannelManager,
|
||||
@@ -317,7 +312,7 @@ func provideMCPManager(log *slog.Logger, service ctr.Service, cfg config.Config,
|
||||
// memory providers
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
func provideMemoryLLM(modelsService *models.Service, queries *dbsqlc.Queries, log *slog.Logger) memory.LLM {
|
||||
func provideMemoryLLM(modelsService *models.Service, queries *dbsqlc.Queries, log *slog.Logger) memprovider.LLM {
|
||||
return &lazyLLMClient{
|
||||
modelsService: modelsService,
|
||||
queries: queries,
|
||||
@@ -326,59 +321,14 @@ func provideMemoryLLM(modelsService *models.Service, queries *dbsqlc.Queries, lo
|
||||
}
|
||||
}
|
||||
|
||||
func provideEmbeddingsResolver(log *slog.Logger, modelsService *models.Service, queries *dbsqlc.Queries) *embeddings.Resolver {
|
||||
return embeddings.NewResolver(log, modelsService, queries, 10*time.Second)
|
||||
}
|
||||
|
||||
type embeddingSetup struct {
|
||||
Vectors map[string]int
|
||||
TextModel models.GetResponse
|
||||
MultimodalModel models.GetResponse
|
||||
HasEmbeddingModels bool
|
||||
}
|
||||
|
||||
func provideEmbeddingSetup(log *slog.Logger, modelsService *models.Service) (embeddingSetup, error) {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
||||
defer cancel()
|
||||
|
||||
vectors, textModel, multimodalModel, hasEmbeddingModels, err := embeddings.CollectEmbeddingVectors(ctx, modelsService)
|
||||
if err != nil {
|
||||
return embeddingSetup{}, fmt.Errorf("embedding models: %w", err)
|
||||
}
|
||||
if hasEmbeddingModels && multimodalModel.ModelID == "" {
|
||||
log.Warn("No multimodal embedding model configured. Multimodal embedding features will be limited.")
|
||||
}
|
||||
return embeddingSetup{
|
||||
Vectors: vectors,
|
||||
TextModel: textModel,
|
||||
MultimodalModel: multimodalModel,
|
||||
HasEmbeddingModels: hasEmbeddingModels,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func provideTextEmbedderForMemory(resolver *embeddings.Resolver, setup embeddingSetup, log *slog.Logger) embeddings.Embedder {
|
||||
return buildTextEmbedder(resolver, setup.TextModel, setup.HasEmbeddingModels, log)
|
||||
}
|
||||
|
||||
func provideQdrantStore(log *slog.Logger, cfg config.Config, setup embeddingSetup) (*memory.QdrantStore, error) {
|
||||
qcfg := cfg.Qdrant
|
||||
timeout := time.Duration(qcfg.TimeoutSeconds) * time.Second
|
||||
if setup.HasEmbeddingModels && len(setup.Vectors) > 0 {
|
||||
store, err := memory.NewQdrantStoreWithVectors(log, qcfg.BaseURL, qcfg.APIKey, qcfg.Collection, setup.Vectors, "sparse_hash", timeout)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("qdrant named vectors init: %w", err)
|
||||
}
|
||||
return store, nil
|
||||
}
|
||||
store, err := memory.NewQdrantStore(log, qcfg.BaseURL, qcfg.APIKey, qcfg.Collection, setup.TextModel.Dimensions, "sparse_hash", timeout)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("qdrant init: %w", err)
|
||||
}
|
||||
return store, nil
|
||||
}
|
||||
|
||||
func provideMemoryService(log *slog.Logger, llm memory.LLM, embedder embeddings.Embedder, store *memory.QdrantStore, resolver *embeddings.Resolver, bm25 *memory.BM25Indexer, setup embeddingSetup) *memory.Service {
|
||||
return memory.NewService(log, llm, embedder, store, resolver, bm25, setup.TextModel.ModelID, setup.MultimodalModel.ModelID)
|
||||
func provideMemoryProviderRegistry(log *slog.Logger, chatService *conversation.Service, accountService *accounts.Service, containerdHandler *handlers.ContainerdHandler) *memprovider.Registry {
|
||||
registry := memprovider.NewRegistry(log)
|
||||
builtinRuntime := handlers.NewBuiltinMemoryRuntime(containerdHandler.FSService())
|
||||
registry.RegisterFactory(memprovider.BuiltinType, func(id string, config map[string]any) (memprovider.Provider, error) {
|
||||
return memprovider.NewBuiltinProvider(log, builtinRuntime, chatService, accountService), nil
|
||||
})
|
||||
registry.Register("__builtin_default__", memprovider.NewBuiltinProvider(log, builtinRuntime, chatService, accountService))
|
||||
return registry
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -405,8 +355,9 @@ func provideHeartbeatTriggerer(resolver *flow.Resolver) heartbeat.Triggerer {
|
||||
// conversation flow
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
func provideChatResolver(log *slog.Logger, cfg config.Config, modelsService *models.Service, queries *dbsqlc.Queries, memoryService *memory.Service, chatService *conversation.Service, msgService *message.DBService, settingsService *settings.Service, mediaService *media.Service, containerdHandler *handlers.ContainerdHandler, inboxService *inbox.Service) *flow.Resolver {
|
||||
resolver := flow.NewResolver(log, modelsService, queries, memoryService, chatService, msgService, settingsService, cfg.AgentGateway.BaseURL(), 120*time.Second)
|
||||
func provideChatResolver(log *slog.Logger, cfg config.Config, modelsService *models.Service, queries *dbsqlc.Queries, chatService *conversation.Service, msgService *message.DBService, settingsService *settings.Service, mediaService *media.Service, containerdHandler *handlers.ContainerdHandler, inboxService *inbox.Service, memoryRegistry *memprovider.Registry) *flow.Resolver {
|
||||
resolver := flow.NewResolver(log, modelsService, queries, chatService, msgService, settingsService, cfg.AgentGateway.BaseURL(), 120*time.Second)
|
||||
resolver.SetMemoryRegistry(memoryRegistry)
|
||||
resolver.SetSkillLoader(&skillLoaderAdapter{handler: containerdHandler})
|
||||
resolver.SetGatewayAssetLoader(&gatewayAssetLoaderAdapter{media: mediaService})
|
||||
resolver.SetInboxService(inboxService)
|
||||
@@ -479,7 +430,7 @@ func provideContainerdHandler(log *slog.Logger, service ctr.Service, manager *mc
|
||||
return handlers.NewContainerdHandler(log, service, manager, cfg.MCP, cfg.Containerd.Namespace, rc.ContainerBackend, botService, accountService, policyService, queries)
|
||||
}
|
||||
|
||||
func provideToolGatewayService(log *slog.Logger, cfg config.Config, channelManager *channel.Manager, registry *channel.Registry, routeService *route.DBService, scheduleService *schedule.Service, memoryService *memory.Service, chatService *conversation.Service, accountService *accounts.Service, settingsService *settings.Service, searchProviderService *searchproviders.Service, manager *mcp.Manager, containerdHandler *handlers.ContainerdHandler, mcpConnService *mcp.ConnectionService, mediaService *media.Service, inboxService *inbox.Service, emailService *emailpkg.Service, emailManager *emailpkg.Manager) *mcp.ToolGatewayService {
|
||||
func provideToolGatewayService(log *slog.Logger, cfg config.Config, channelManager *channel.Manager, registry *channel.Registry, routeService *route.DBService, scheduleService *schedule.Service, chatService *conversation.Service, accountService *accounts.Service, settingsService *settings.Service, searchProviderService *searchproviders.Service, manager *mcp.Manager, containerdHandler *handlers.ContainerdHandler, mcpConnService *mcp.ConnectionService, mediaService *media.Service, inboxService *inbox.Service, memoryRegistry *memprovider.Registry, emailService *emailpkg.Service, emailManager *emailpkg.Manager) *mcp.ToolGatewayService {
|
||||
var assetResolver mcpmessage.AssetResolver
|
||||
if mediaService != nil {
|
||||
assetResolver = &mediaAssetResolverAdapter{media: mediaService}
|
||||
@@ -487,7 +438,7 @@ func provideToolGatewayService(log *slog.Logger, cfg config.Config, channelManag
|
||||
messageExec := mcpmessage.NewExecutor(log, channelManager, channelManager, registry, assetResolver)
|
||||
contactsExec := mcpcontacts.NewExecutor(log, routeService)
|
||||
scheduleExec := mcpschedule.NewExecutor(log, scheduleService)
|
||||
memoryExec := mcpmemory.NewExecutor(log, memoryService, chatService, accountService)
|
||||
memoryExec := mcpmemory.NewExecutor(log, memoryRegistry, settingsService)
|
||||
webExec := mcpweb.NewExecutor(log, settingsService, searchProviderService)
|
||||
inboxExec := mcpinbox.NewExecutor(log, inboxService)
|
||||
fsExec := mcpcontainer.NewExecutor(log, manager, config.DefaultDataMount)
|
||||
@@ -510,12 +461,11 @@ func provideToolGatewayService(log *slog.Logger, cfg config.Config, channelManag
|
||||
// handler providers (interface adaptation / config extraction)
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
func provideMemoryHandler(log *slog.Logger, service *memory.Service, chatService *conversation.Service, accountService *accounts.Service, cfg config.Config, manager *mcp.Manager) *handlers.MemoryHandler {
|
||||
h := handlers.NewMemoryHandler(log, service, chatService, accountService)
|
||||
if manager != nil {
|
||||
execWorkDir := config.DefaultDataMount
|
||||
h.SetMemoryFS(memory.NewMemoryFS(log, manager, execWorkDir))
|
||||
}
|
||||
func provideMemoryHandler(log *slog.Logger, botService *bots.Service, accountService *accounts.Service, cfg config.Config, manager *mcp.Manager, memoryRegistry *memprovider.Registry, settingsService *settings.Service, containerdHandler *handlers.ContainerdHandler) *handlers.MemoryHandler {
|
||||
h := handlers.NewMemoryHandler(log, botService, accountService)
|
||||
h.SetMemoryRegistry(memoryRegistry)
|
||||
h.SetSettingsService(settingsService)
|
||||
h.SetFSService(containerdHandler.FSService())
|
||||
return h
|
||||
}
|
||||
|
||||
@@ -615,14 +565,19 @@ func provideServer(params serverParams) *server.Server {
|
||||
// lifecycle hooks
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
func startMemoryWarmup(lc fx.Lifecycle, memoryService *memory.Service, logger *slog.Logger) {
|
||||
func startMemoryProviderBootstrap(lc fx.Lifecycle, log *slog.Logger, mpService *memprovider.Service, registry *memprovider.Registry) {
|
||||
lc.Append(fx.Hook{
|
||||
OnStart: func(ctx context.Context) error {
|
||||
go func() {
|
||||
if err := memoryService.WarmupBM25(context.Background(), 200); err != nil {
|
||||
logger.Warn("bm25 warmup failed", slog.Any("error", err))
|
||||
}
|
||||
}()
|
||||
resp, err := mpService.EnsureDefault(ctx)
|
||||
if err != nil {
|
||||
log.Warn("failed to ensure default memory provider", slog.Any("error", err))
|
||||
return nil
|
||||
}
|
||||
if _, regErr := registry.Instantiate(resp.ID, resp.Provider, resp.Config); regErr != nil {
|
||||
log.Warn("failed to instantiate default memory provider", slog.Any("error", regErr))
|
||||
} else {
|
||||
log.Info("default memory provider ready", slog.String("id", resp.ID), slog.String("provider", resp.Provider))
|
||||
}
|
||||
return nil
|
||||
},
|
||||
})
|
||||
@@ -707,21 +662,6 @@ func startServer(lc fx.Lifecycle, logger *slog.Logger, srv *server.Server, shutd
|
||||
// helpers
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
func buildTextEmbedder(resolver *embeddings.Resolver, textModel models.GetResponse, hasModels bool, log *slog.Logger) embeddings.Embedder {
|
||||
if !hasModels {
|
||||
return nil
|
||||
}
|
||||
if textModel.ModelID == "" || textModel.Dimensions <= 0 {
|
||||
log.Warn("No text embedding model configured. Text embedding features will be limited.")
|
||||
return nil
|
||||
}
|
||||
return &embeddings.ResolverTextEmbedder{
|
||||
Resolver: resolver,
|
||||
ModelID: textModel.ModelID,
|
||||
Dims: textModel.Dimensions,
|
||||
}
|
||||
}
|
||||
|
||||
func ensureAdminUser(ctx context.Context, log *slog.Logger, queries *dbsqlc.Queries, cfg config.Config) error {
|
||||
if queries == nil {
|
||||
return fmt.Errorf("db queries not configured")
|
||||
@@ -793,26 +733,26 @@ type lazyLLMClient struct {
|
||||
logger *slog.Logger
|
||||
}
|
||||
|
||||
func (c *lazyLLMClient) Extract(ctx context.Context, req memory.ExtractRequest) (memory.ExtractResponse, error) {
|
||||
func (c *lazyLLMClient) Extract(ctx context.Context, req memprovider.ExtractRequest) (memprovider.ExtractResponse, error) {
|
||||
client, err := c.resolve(ctx)
|
||||
if err != nil {
|
||||
return memory.ExtractResponse{}, err
|
||||
return memprovider.ExtractResponse{}, err
|
||||
}
|
||||
return client.Extract(ctx, req)
|
||||
}
|
||||
|
||||
func (c *lazyLLMClient) Decide(ctx context.Context, req memory.DecideRequest) (memory.DecideResponse, error) {
|
||||
func (c *lazyLLMClient) Decide(ctx context.Context, req memprovider.DecideRequest) (memprovider.DecideResponse, error) {
|
||||
client, err := c.resolve(ctx)
|
||||
if err != nil {
|
||||
return memory.DecideResponse{}, err
|
||||
return memprovider.DecideResponse{}, err
|
||||
}
|
||||
return client.Decide(ctx, req)
|
||||
}
|
||||
|
||||
func (c *lazyLLMClient) Compact(ctx context.Context, req memory.CompactRequest) (memory.CompactResponse, error) {
|
||||
func (c *lazyLLMClient) Compact(ctx context.Context, req memprovider.CompactRequest) (memprovider.CompactResponse, error) {
|
||||
client, err := c.resolve(ctx)
|
||||
if err != nil {
|
||||
return memory.CompactResponse{}, err
|
||||
return memprovider.CompactResponse{}, err
|
||||
}
|
||||
return client.Compact(ctx, req)
|
||||
}
|
||||
@@ -825,11 +765,11 @@ func (c *lazyLLMClient) DetectLanguage(ctx context.Context, text string) (string
|
||||
return client.DetectLanguage(ctx, text)
|
||||
}
|
||||
|
||||
func (c *lazyLLMClient) resolve(ctx context.Context) (memory.LLM, error) {
|
||||
func (c *lazyLLMClient) resolve(ctx context.Context) (memprovider.LLM, error) {
|
||||
if c.modelsService == nil || c.queries == nil {
|
||||
return nil, fmt.Errorf("models service not configured")
|
||||
}
|
||||
botID := memory.BotIDFromContext(ctx)
|
||||
botID := ""
|
||||
memoryModel, memoryProvider, err := models.SelectMemoryModelForBot(ctx, c.modelsService, c.queries, botID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -840,7 +780,9 @@ func (c *lazyLLMClient) resolve(ctx context.Context) (memory.LLM, error) {
|
||||
default:
|
||||
return nil, fmt.Errorf("memory model client type not supported: %s", clientType)
|
||||
}
|
||||
return memory.NewLLMClient(c.logger, memoryProvider.BaseUrl, memoryProvider.ApiKey, memoryModel.ModelID, c.timeout)
|
||||
_ = memoryProvider
|
||||
_ = memoryModel
|
||||
return nil, fmt.Errorf("memory llm runtime is not available")
|
||||
}
|
||||
|
||||
// skillLoaderAdapter bridges handlers.ContainerdHandler to flow.SkillLoader.
|
||||
|
||||
Reference in New Issue
Block a user