feat(020): US2 — per-agent core memory blob
Letta-style identity blob, one per (owner, agent), always included in
session-start (my_status) responses. Replaces wholesale on rewrite; size
capped at SYNAPBUS_CORE_MEMORY_MAX_BYTES (default 2048); owner-scoped.
Components:
- internal/messaging/memory_core.go (+ test): CoreMemoryStore with
Get/Set/Delete/List, ErrCoreMemoryTooLarge, NewCoreProvider adapter
for search.CoreMemoryProvider.
- internal/mcp/server.go SetInjection: wires the core provider into
the my_status handler wrap.
- internal/mcp/injection_core_test.go: seed → wrapped my_status →
relevant_context.core_memory matches; missing row → no field.
- internal/api/memory_core.go + router: GET/PUT/DELETE
/api/owner/{ownerID}/agents/{agentName}/core-memory, session-auth,
413 on oversize.
- internal/admin/socket.go: memory.core.{get,set,delete} dispatch
handlers with username→user.id resolution.
- cmd/synapbus/admin.go: `synapbus memory core {get,set,delete}` cobra
subtree.
- cmd/synapbus/main.go: wires ParseMemoryConfig, CoreMemoryStore,
MemoryInjections; calls mcpSrv.SetInjection on startup.
Pin overlay still TODO (US3-T029).
Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 4.7
parent
8a5d5e1f59
commit
a52d68ed88
+103
-5
@@ -702,10 +702,10 @@ func addAdminCommands(rootCmd *cobra.Command) {
|
||||
channelsJoinCmd.MarkFlagRequired("agent")
|
||||
|
||||
var (
|
||||
channelsUpdateName string
|
||||
channelsUpdateAutoApprove string
|
||||
channelsUpdateStalemateRemind string
|
||||
channelsUpdateStalemateEscalate string
|
||||
channelsUpdateName string
|
||||
channelsUpdateAutoApprove string
|
||||
channelsUpdateStalemateRemind string
|
||||
channelsUpdateStalemateEscalate string
|
||||
)
|
||||
channelsUpdateCmd := &cobra.Command{
|
||||
Use: "update",
|
||||
@@ -1341,10 +1341,108 @@ Examples:
|
||||
harnessConfigCmd.AddCommand(harnessConfigGetCmd, harnessConfigSetCmd, harnessConfigEditCmd)
|
||||
harnessCmd.AddCommand(harnessConfigCmd)
|
||||
|
||||
// ----- memory commands (feature 020 — proactive memory + dream worker) -----
|
||||
memoryCmd := &cobra.Command{
|
||||
Use: "memory",
|
||||
Short: "Manage proactive memory (feature 020)",
|
||||
}
|
||||
|
||||
memoryCoreCmd := &cobra.Command{
|
||||
Use: "core",
|
||||
Short: "Manage per-(owner, agent) core memory blobs",
|
||||
}
|
||||
|
||||
var (
|
||||
memCoreOwner string
|
||||
memCoreAgent string
|
||||
memCoreBlob string
|
||||
memCoreBlobFile string
|
||||
memCoreUpdater string
|
||||
)
|
||||
|
||||
memoryCoreGetCmd := &cobra.Command{
|
||||
Use: "get",
|
||||
Short: "Print the current core memory blob for an (owner, agent)",
|
||||
RunE: func(cmd *cobra.Command, args []string) error {
|
||||
resp, err := adminRequest("memory.core.get", map[string]string{
|
||||
"owner": memCoreOwner,
|
||||
"agent": memCoreAgent,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
printJSON(resp["data"])
|
||||
return nil
|
||||
},
|
||||
}
|
||||
memoryCoreGetCmd.Flags().StringVar(&memCoreOwner, "owner", "", "Owner username or numeric user ID")
|
||||
memoryCoreGetCmd.Flags().StringVar(&memCoreAgent, "agent", "", "Agent name")
|
||||
_ = memoryCoreGetCmd.MarkFlagRequired("owner")
|
||||
_ = memoryCoreGetCmd.MarkFlagRequired("agent")
|
||||
|
||||
memoryCoreSetCmd := &cobra.Command{
|
||||
Use: "set",
|
||||
Short: "Replace the core memory blob (wholesale, no merge)",
|
||||
RunE: func(cmd *cobra.Command, args []string) error {
|
||||
blob := memCoreBlob
|
||||
if memCoreBlobFile != "" {
|
||||
data, err := os.ReadFile(memCoreBlobFile)
|
||||
if err != nil {
|
||||
return fmt.Errorf("read --blob-file: %w", err)
|
||||
}
|
||||
blob = string(data)
|
||||
}
|
||||
if blob == "" {
|
||||
return fmt.Errorf("either --blob or --blob-file is required (non-empty)")
|
||||
}
|
||||
resp, err := adminRequest("memory.core.set", map[string]string{
|
||||
"owner": memCoreOwner,
|
||||
"agent": memCoreAgent,
|
||||
"blob": blob,
|
||||
"updated_by": memCoreUpdater,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
printJSON(resp["data"])
|
||||
return nil
|
||||
},
|
||||
}
|
||||
memoryCoreSetCmd.Flags().StringVar(&memCoreOwner, "owner", "", "Owner username or numeric user ID")
|
||||
memoryCoreSetCmd.Flags().StringVar(&memCoreAgent, "agent", "", "Agent name")
|
||||
memoryCoreSetCmd.Flags().StringVar(&memCoreBlob, "blob", "", "Core memory blob (inline)")
|
||||
memoryCoreSetCmd.Flags().StringVar(&memCoreBlobFile, "blob-file", "", "Read blob from file path (overrides --blob)")
|
||||
memoryCoreSetCmd.Flags().StringVar(&memCoreUpdater, "updated-by", "human", "updated_by audit field (default: human)")
|
||||
_ = memoryCoreSetCmd.MarkFlagRequired("owner")
|
||||
_ = memoryCoreSetCmd.MarkFlagRequired("agent")
|
||||
|
||||
memoryCoreDeleteCmd := &cobra.Command{
|
||||
Use: "delete",
|
||||
Short: "Remove the core memory blob",
|
||||
RunE: func(cmd *cobra.Command, args []string) error {
|
||||
resp, err := adminRequest("memory.core.delete", map[string]string{
|
||||
"owner": memCoreOwner,
|
||||
"agent": memCoreAgent,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
printJSON(resp["data"])
|
||||
return nil
|
||||
},
|
||||
}
|
||||
memoryCoreDeleteCmd.Flags().StringVar(&memCoreOwner, "owner", "", "Owner username or numeric user ID")
|
||||
memoryCoreDeleteCmd.Flags().StringVar(&memCoreAgent, "agent", "", "Agent name")
|
||||
_ = memoryCoreDeleteCmd.MarkFlagRequired("owner")
|
||||
_ = memoryCoreDeleteCmd.MarkFlagRequired("agent")
|
||||
|
||||
memoryCoreCmd.AddCommand(memoryCoreGetCmd, memoryCoreSetCmd, memoryCoreDeleteCmd)
|
||||
memoryCmd.AddCommand(memoryCoreCmd)
|
||||
|
||||
// ----- add persistent flag and commands to root -----
|
||||
rootCmd.PersistentFlags().StringVar(&adminSocket, "socket", "/tmp/synapbus.sock", "Path to admin Unix socket")
|
||||
|
||||
rootCmd.AddCommand(userCmd, agentCmd, auditCmd, backupCmd, messagesCmd, channelsCmd, conversationsCmd, embeddingsCmd, dbCmd, retentionCmd, webhookCmd, k8sCmd, attachmentsCmd, harnessCmd)
|
||||
rootCmd.AddCommand(userCmd, agentCmd, auditCmd, backupCmd, messagesCmd, channelsCmd, conversationsCmd, embeddingsCmd, dbCmd, retentionCmd, webhookCmd, k8sCmd, attachmentsCmd, harnessCmd, memoryCmd)
|
||||
}
|
||||
|
||||
// toTableRows remaps []map[string]string using a header->key mapping.
|
||||
|
||||
+49
-29
@@ -26,6 +26,7 @@ import (
|
||||
"github.com/synapbus/synapbus/internal/a2a"
|
||||
"github.com/synapbus/synapbus/internal/actions"
|
||||
"github.com/synapbus/synapbus/internal/admin"
|
||||
"github.com/synapbus/synapbus/internal/agentquery"
|
||||
"github.com/synapbus/synapbus/internal/agents"
|
||||
"github.com/synapbus/synapbus/internal/api"
|
||||
"github.com/synapbus/synapbus/internal/apikeys"
|
||||
@@ -34,51 +35,50 @@ import (
|
||||
"github.com/synapbus/synapbus/internal/auth/idp"
|
||||
"github.com/synapbus/synapbus/internal/channels"
|
||||
"github.com/synapbus/synapbus/internal/console"
|
||||
"github.com/synapbus/synapbus/internal/dispatcher"
|
||||
"github.com/synapbus/synapbus/internal/goals"
|
||||
"github.com/synapbus/synapbus/internal/goaltasks"
|
||||
"github.com/synapbus/synapbus/internal/secrets"
|
||||
"github.com/synapbus/synapbus/internal/dispatcher"
|
||||
"github.com/synapbus/synapbus/internal/health"
|
||||
"github.com/synapbus/synapbus/internal/jsruntime"
|
||||
k8spkg "github.com/synapbus/synapbus/internal/k8s"
|
||||
"github.com/synapbus/synapbus/internal/marketplace"
|
||||
mcpserver "github.com/synapbus/synapbus/internal/mcp"
|
||||
"github.com/synapbus/synapbus/internal/agentquery"
|
||||
reactorpkg "github.com/synapbus/synapbus/internal/reactor"
|
||||
"github.com/synapbus/synapbus/internal/messaging"
|
||||
prommetrics "github.com/synapbus/synapbus/internal/metrics"
|
||||
"github.com/synapbus/synapbus/internal/harness"
|
||||
"github.com/synapbus/synapbus/internal/harness/docker"
|
||||
"github.com/synapbus/synapbus/internal/harness/k8sjob"
|
||||
"github.com/synapbus/synapbus/internal/harness/runs"
|
||||
"github.com/synapbus/synapbus/internal/harness/subprocess"
|
||||
"github.com/synapbus/synapbus/internal/harness/webhook"
|
||||
"github.com/synapbus/synapbus/internal/health"
|
||||
"github.com/synapbus/synapbus/internal/jsruntime"
|
||||
k8spkg "github.com/synapbus/synapbus/internal/k8s"
|
||||
"github.com/synapbus/synapbus/internal/marketplace"
|
||||
mcpserver "github.com/synapbus/synapbus/internal/mcp"
|
||||
"github.com/synapbus/synapbus/internal/messaging"
|
||||
prommetrics "github.com/synapbus/synapbus/internal/metrics"
|
||||
"github.com/synapbus/synapbus/internal/observability"
|
||||
"github.com/synapbus/synapbus/internal/push"
|
||||
"github.com/synapbus/synapbus/internal/reactions"
|
||||
reactorpkg "github.com/synapbus/synapbus/internal/reactor"
|
||||
"github.com/synapbus/synapbus/internal/search"
|
||||
"github.com/synapbus/synapbus/internal/search/embedding"
|
||||
"github.com/synapbus/synapbus/internal/secrets"
|
||||
"github.com/synapbus/synapbus/internal/storage"
|
||||
"github.com/synapbus/synapbus/internal/push"
|
||||
"github.com/synapbus/synapbus/internal/trace"
|
||||
"github.com/synapbus/synapbus/internal/trust"
|
||||
"github.com/synapbus/synapbus/internal/web"
|
||||
"github.com/synapbus/synapbus/internal/wiki"
|
||||
"github.com/synapbus/synapbus/internal/webhooks"
|
||||
"github.com/synapbus/synapbus/internal/wiki"
|
||||
)
|
||||
|
||||
// version is set at build time via -ldflags "-X main.version=..."
|
||||
var version = "dev"
|
||||
|
||||
var (
|
||||
host string
|
||||
port int
|
||||
dataDir string
|
||||
logLevel string
|
||||
metricsEnabled bool
|
||||
traceRetention string
|
||||
adminSocketPath string
|
||||
webhookWorkers int
|
||||
messageRetention string
|
||||
host string
|
||||
port int
|
||||
dataDir string
|
||||
logLevel string
|
||||
metricsEnabled bool
|
||||
traceRetention string
|
||||
adminSocketPath string
|
||||
webhookWorkers int
|
||||
messageRetention string
|
||||
)
|
||||
|
||||
func main() {
|
||||
@@ -582,6 +582,24 @@ func runServe(cmd *cobra.Command, args []string) error {
|
||||
|
||||
mcpSrv := mcpserver.NewMCPServer(msgService, agentService, channelService, swarmService, attachmentService, searchService, reactionService, trustService, wikiService, con, jsPool, actionRegistry, actionIndex, db.DB)
|
||||
|
||||
// Feature 020 — proactive memory injection.
|
||||
//
|
||||
// Parse the memory config from env, build the audit-ring store and
|
||||
// the per-(owner, agent) core memory store, and wire both into the
|
||||
// MCP hybrid tool surface. When SYNAPBUS_INJECTION_ENABLED=0 (the
|
||||
// default), SetInjection still runs but WrapInjection returns each
|
||||
// handler unchanged, so tool responses keep their pre-feature shape
|
||||
// bit-for-bit (FR-012, SC-009).
|
||||
memCfg := messaging.ParseMemoryConfig()
|
||||
memoryInjectionStore := messaging.NewMemoryInjections(db.DB)
|
||||
coreMemoryStore := messaging.NewCoreMemoryStore(db.DB, memCfg.CoreMemoryMaxBytes)
|
||||
mcpSrv.SetInjection(memCfg, memoryInjectionStore, messaging.NewCoreProvider(coreMemoryStore))
|
||||
slog.Info("proactive memory injection wired",
|
||||
"enabled", memCfg.InjectionEnabled,
|
||||
"budget_tokens", memCfg.InjectionBudgetTokens,
|
||||
"core_max_bytes", memCfg.CoreMemoryMaxBytes,
|
||||
)
|
||||
|
||||
// Wire the agent marketplace (spec 016 MVP).
|
||||
marketplaceStore := marketplace.NewStore(db.DB)
|
||||
marketplaceSvc := marketplace.NewService(marketplaceStore, wikiService, swarmService, channelService, msgService, tracer)
|
||||
@@ -789,6 +807,7 @@ func runServe(cmd *cobra.Command, args []string) error {
|
||||
GoalTasksService: goalTasksService,
|
||||
BaseURL: baseURL,
|
||||
WikiService: wikiService,
|
||||
CoreMemoryStore: coreMemoryStore,
|
||||
})
|
||||
r.Mount("/", apiRouter)
|
||||
|
||||
@@ -797,13 +816,13 @@ func runServe(cmd *cobra.Command, args []string) error {
|
||||
|
||||
// Start admin socket server
|
||||
adminSvcs := &admin.Services{
|
||||
Users: userStore,
|
||||
Sessions: sessionStore,
|
||||
Agents: agentService,
|
||||
Messages: msgService,
|
||||
Channels: channelService,
|
||||
Traces: traceStore,
|
||||
DataDir: dataDir,
|
||||
Users: userStore,
|
||||
Sessions: sessionStore,
|
||||
Agents: agentService,
|
||||
Messages: msgService,
|
||||
Channels: channelService,
|
||||
Traces: traceStore,
|
||||
DataDir: dataDir,
|
||||
}
|
||||
// Wire optional services into admin (may be nil if not configured)
|
||||
if searchCfg.IsEnabled() {
|
||||
@@ -819,6 +838,7 @@ func runServe(cmd *cobra.Command, args []string) error {
|
||||
}
|
||||
adminSvcs.WebhookService = webhookService
|
||||
adminSvcs.K8sService = k8sService
|
||||
adminSvcs.CoreMemoryStore = coreMemoryStore
|
||||
adminServer := admin.NewServer(adminSocketPath, db.DB, adminSvcs, logger)
|
||||
if err := adminServer.Start(); err != nil {
|
||||
return fmt.Errorf("start admin socket: %w", err)
|
||||
|
||||
+19
-14
@@ -34,20 +34,25 @@ type K8sServiceProvider interface {
|
||||
|
||||
// Services holds references to all services the admin socket can control.
|
||||
type Services struct {
|
||||
Users *auth.SQLiteUserStore
|
||||
Sessions auth.SessionStore
|
||||
Agents *agents.AgentService
|
||||
Messages *messaging.MessagingService
|
||||
Channels *channels.Service
|
||||
Traces trace.TraceStore
|
||||
EmbeddingStore *search.EmbeddingStore
|
||||
VectorIndex *search.VectorIndex
|
||||
SearchService *search.Service
|
||||
AttachmentService *attachments.Service
|
||||
WebhookService WebhookServiceProvider
|
||||
K8sService K8sServiceProvider
|
||||
DataDir string
|
||||
RetentionWorker RetentionStatusProvider
|
||||
Users *auth.SQLiteUserStore
|
||||
Sessions auth.SessionStore
|
||||
Agents *agents.AgentService
|
||||
Messages *messaging.MessagingService
|
||||
Channels *channels.Service
|
||||
Traces trace.TraceStore
|
||||
EmbeddingStore *search.EmbeddingStore
|
||||
VectorIndex *search.VectorIndex
|
||||
SearchService *search.Service
|
||||
AttachmentService *attachments.Service
|
||||
WebhookService WebhookServiceProvider
|
||||
K8sService K8sServiceProvider
|
||||
DataDir string
|
||||
RetentionWorker RetentionStatusProvider
|
||||
|
||||
// CoreMemoryStore is the per-(owner, agent) core memory store wired
|
||||
// in for feature 020 admin CLI commands (`synapbus memory core ...`).
|
||||
// May be nil — handlers report "core memory store not configured".
|
||||
CoreMemoryStore *messaging.CoreMemoryStore
|
||||
}
|
||||
|
||||
// RetentionStatusProvider provides retention status information.
|
||||
|
||||
+137
-2
@@ -226,6 +226,14 @@ func (s *AdminServer) dispatch(req Request) Response {
|
||||
case "harness.config_set":
|
||||
return s.handleHarnessConfigSet(ctx, req.Args)
|
||||
|
||||
// --- memory core (feature 020 — proactive memory) ---
|
||||
case "memory.core.get":
|
||||
return s.handleMemoryCoreGet(ctx, req.Args)
|
||||
case "memory.core.set":
|
||||
return s.handleMemoryCoreSet(ctx, req.Args)
|
||||
case "memory.core.delete":
|
||||
return s.handleMemoryCoreDelete(ctx, req.Args)
|
||||
|
||||
default:
|
||||
return Response{OK: false, Error: fmt.Sprintf("unknown command: %s", req.Command)}
|
||||
}
|
||||
@@ -545,7 +553,7 @@ func (s *AdminServer) handleAuditStats(ctx context.Context) Response {
|
||||
}
|
||||
|
||||
return Response{OK: true, Data: map[string]interface{}{
|
||||
"total_traces": totalTraces,
|
||||
"total_traces": totalTraces,
|
||||
"counts_by_action": counts,
|
||||
}}
|
||||
}
|
||||
@@ -1694,7 +1702,7 @@ func (s *AdminServer) handleAttachmentsGC(ctx context.Context) Response {
|
||||
}
|
||||
|
||||
return Response{OK: true, Data: map[string]interface{}{
|
||||
"files_removed": result.FilesRemoved,
|
||||
"files_removed": result.FilesRemoved,
|
||||
"bytes_reclaimed": result.BytesReclaimed,
|
||||
}}
|
||||
}
|
||||
@@ -1850,5 +1858,132 @@ func (s *AdminServer) handleHarnessConfigSet(ctx context.Context, args json.RawM
|
||||
}}
|
||||
}
|
||||
|
||||
// ---------- memory core handlers (feature 020) ----------
|
||||
//
|
||||
// owner_id wire format: callers pass the owner as a username; we resolve
|
||||
// to `users.id` and pass the string form to CoreMemoryStore so it matches
|
||||
// the proactive-memory tables' TEXT owner_id convention.
|
||||
|
||||
func (s *AdminServer) resolveOwnerString(ctx context.Context, ownerInput string) (string, error) {
|
||||
if ownerInput == "" {
|
||||
return "", fmt.Errorf("owner is required")
|
||||
}
|
||||
// First try numeric — admins may already know the user ID.
|
||||
var id int64
|
||||
if _, err := fmt.Sscanf(ownerInput, "%d", &id); err == nil && id > 0 {
|
||||
return fmt.Sprintf("%d", id), nil
|
||||
}
|
||||
// Otherwise treat as username.
|
||||
user, err := s.services.Users.GetUserByUsername(ctx, ownerInput)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("resolve owner %q: %w", ownerInput, err)
|
||||
}
|
||||
return fmt.Sprintf("%d", user.ID), nil
|
||||
}
|
||||
|
||||
func (s *AdminServer) handleMemoryCoreGet(ctx context.Context, args json.RawMessage) Response {
|
||||
var p struct {
|
||||
Owner string `json:"owner"`
|
||||
Agent string `json:"agent"`
|
||||
}
|
||||
if err := json.Unmarshal(args, &p); err != nil {
|
||||
return Response{OK: false, Error: "invalid args: " + err.Error()}
|
||||
}
|
||||
if p.Agent == "" {
|
||||
return Response{OK: false, Error: "agent is required"}
|
||||
}
|
||||
if s.services.CoreMemoryStore == nil {
|
||||
return Response{OK: false, Error: "core memory store not configured"}
|
||||
}
|
||||
ownerStr, err := s.resolveOwnerString(ctx, p.Owner)
|
||||
if err != nil {
|
||||
return Response{OK: false, Error: err.Error()}
|
||||
}
|
||||
blob, updatedAt, ok, err := s.services.CoreMemoryStore.Get(ctx, ownerStr, p.Agent)
|
||||
if err != nil {
|
||||
return Response{OK: false, Error: err.Error()}
|
||||
}
|
||||
if !ok {
|
||||
return Response{OK: true, Data: map[string]any{
|
||||
"owner_id": ownerStr,
|
||||
"agent_name": p.Agent,
|
||||
"exists": false,
|
||||
}}
|
||||
}
|
||||
return Response{OK: true, Data: map[string]any{
|
||||
"owner_id": ownerStr,
|
||||
"agent_name": p.Agent,
|
||||
"exists": true,
|
||||
"blob": blob,
|
||||
"updated_at": updatedAt.Format(time.RFC3339),
|
||||
}}
|
||||
}
|
||||
|
||||
func (s *AdminServer) handleMemoryCoreSet(ctx context.Context, args json.RawMessage) Response {
|
||||
var p struct {
|
||||
Owner string `json:"owner"`
|
||||
Agent string `json:"agent"`
|
||||
Blob string `json:"blob"`
|
||||
UpdatedBy string `json:"updated_by"`
|
||||
}
|
||||
if err := json.Unmarshal(args, &p); err != nil {
|
||||
return Response{OK: false, Error: "invalid args: " + err.Error()}
|
||||
}
|
||||
if p.Agent == "" {
|
||||
return Response{OK: false, Error: "agent is required"}
|
||||
}
|
||||
if s.services.CoreMemoryStore == nil {
|
||||
return Response{OK: false, Error: "core memory store not configured"}
|
||||
}
|
||||
ownerStr, err := s.resolveOwnerString(ctx, p.Owner)
|
||||
if err != nil {
|
||||
return Response{OK: false, Error: err.Error()}
|
||||
}
|
||||
updatedBy := p.UpdatedBy
|
||||
if updatedBy == "" {
|
||||
updatedBy = "human"
|
||||
}
|
||||
if err := s.services.CoreMemoryStore.Set(ctx, ownerStr, p.Agent, p.Blob, updatedBy); err != nil {
|
||||
if err == messaging.ErrCoreMemoryTooLarge {
|
||||
return Response{OK: false, Error: fmt.Sprintf("core_memory_too_large: blob exceeds %d bytes", s.services.CoreMemoryStore.MaxBytes())}
|
||||
}
|
||||
return Response{OK: false, Error: err.Error()}
|
||||
}
|
||||
return Response{OK: true, Data: map[string]any{
|
||||
"owner_id": ownerStr,
|
||||
"agent_name": p.Agent,
|
||||
"blob_chars": len(p.Blob),
|
||||
"updated_by": updatedBy,
|
||||
}}
|
||||
}
|
||||
|
||||
func (s *AdminServer) handleMemoryCoreDelete(ctx context.Context, args json.RawMessage) Response {
|
||||
var p struct {
|
||||
Owner string `json:"owner"`
|
||||
Agent string `json:"agent"`
|
||||
}
|
||||
if err := json.Unmarshal(args, &p); err != nil {
|
||||
return Response{OK: false, Error: "invalid args: " + err.Error()}
|
||||
}
|
||||
if p.Agent == "" {
|
||||
return Response{OK: false, Error: "agent is required"}
|
||||
}
|
||||
if s.services.CoreMemoryStore == nil {
|
||||
return Response{OK: false, Error: "core memory store not configured"}
|
||||
}
|
||||
ownerStr, err := s.resolveOwnerString(ctx, p.Owner)
|
||||
if err != nil {
|
||||
return Response{OK: false, Error: err.Error()}
|
||||
}
|
||||
if err := s.services.CoreMemoryStore.Delete(ctx, ownerStr, p.Agent); err != nil {
|
||||
return Response{OK: false, Error: err.Error()}
|
||||
}
|
||||
return Response{OK: true, Data: map[string]any{
|
||||
"owner_id": ownerStr,
|
||||
"agent_name": p.Agent,
|
||||
"deleted": true,
|
||||
}}
|
||||
}
|
||||
|
||||
// Ensure the messaging import is used.
|
||||
var _ = messaging.StatusPending
|
||||
|
||||
@@ -0,0 +1,175 @@
|
||||
// REST endpoints for per-(owner, agent) core memory (feature 020 — US2).
|
||||
// Surfaces the underlying messaging.CoreMemoryStore to the Web UI under
|
||||
// `/api/owner/{ownerID}/agents/{agentName}/core-memory`.
|
||||
//
|
||||
// Auth: every handler enforces that the session-bound owner matches the
|
||||
// path's `ownerID`. Cross-owner access yields 403 to avoid leaking the
|
||||
// existence of another owner's resources.
|
||||
package api
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"strconv"
|
||||
|
||||
"github.com/go-chi/chi/v5"
|
||||
|
||||
"github.com/synapbus/synapbus/internal/messaging"
|
||||
)
|
||||
|
||||
// MemoryCoreHandler exposes GET/PUT/DELETE for the `memory_core` table.
|
||||
type MemoryCoreHandler struct {
|
||||
store *messaging.CoreMemoryStore
|
||||
logger *slog.Logger
|
||||
}
|
||||
|
||||
// NewMemoryCoreHandler wires the handler. `store` must be non-nil.
|
||||
func NewMemoryCoreHandler(store *messaging.CoreMemoryStore) *MemoryCoreHandler {
|
||||
return &MemoryCoreHandler{
|
||||
store: store,
|
||||
logger: slog.Default().With("component", "api.memory-core"),
|
||||
}
|
||||
}
|
||||
|
||||
// authorize resolves the URL `ownerID` and confirms it matches the
|
||||
// session-bound owner. Returns the resolved owner string (matching the
|
||||
// memory_core.owner_id TEXT format) plus the agent name on success.
|
||||
func (h *MemoryCoreHandler) authorize(w http.ResponseWriter, r *http.Request) (ownerStr string, agentName string, ok bool) {
|
||||
sessionOwnerID, found := OwnerIDFromContext(r.Context())
|
||||
if !found {
|
||||
writeJSON(w, http.StatusUnauthorized, errorBody("unauthorized", "Authentication required"))
|
||||
return "", "", false
|
||||
}
|
||||
|
||||
pathOwner := chi.URLParam(r, "ownerID")
|
||||
if pathOwner == "" {
|
||||
writeJSON(w, http.StatusBadRequest, errorBody("bad_request", "ownerID is required"))
|
||||
return "", "", false
|
||||
}
|
||||
pathOwnerID, err := strconv.ParseInt(pathOwner, 10, 64)
|
||||
if err != nil || pathOwnerID <= 0 {
|
||||
writeJSON(w, http.StatusBadRequest, errorBody("bad_request", "ownerID must be a positive integer"))
|
||||
return "", "", false
|
||||
}
|
||||
if pathOwnerID != sessionOwnerID {
|
||||
// Use 403, not 404, so the response shape matches other owner-scoped
|
||||
// handlers in this package (see agents_handler.GetAgent).
|
||||
writeJSON(w, http.StatusForbidden, errorBody("forbidden", "You do not have access to this owner"))
|
||||
return "", "", false
|
||||
}
|
||||
|
||||
agentName = chi.URLParam(r, "agentName")
|
||||
if agentName == "" {
|
||||
writeJSON(w, http.StatusBadRequest, errorBody("bad_request", "agentName is required"))
|
||||
return "", "", false
|
||||
}
|
||||
|
||||
return strconv.FormatInt(pathOwnerID, 10), agentName, true
|
||||
}
|
||||
|
||||
// Get handles GET /api/owner/{ownerID}/agents/{agentName}/core-memory.
|
||||
func (h *MemoryCoreHandler) Get(w http.ResponseWriter, r *http.Request) {
|
||||
if h.store == nil {
|
||||
writeJSON(w, http.StatusServiceUnavailable, errorBody("unavailable", "core memory store not configured"))
|
||||
return
|
||||
}
|
||||
ownerStr, agentName, ok := h.authorize(w, r)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
blob, updatedAt, exists, err := h.store.Get(r.Context(), ownerStr, agentName)
|
||||
if err != nil {
|
||||
h.logger.Error("memory_core get failed", "error", err, "owner", ownerStr, "agent", agentName)
|
||||
writeJSON(w, http.StatusInternalServerError, errorBody("server_error", "Failed to read core memory"))
|
||||
return
|
||||
}
|
||||
if !exists {
|
||||
writeJSON(w, http.StatusNotFound, errorBody("not_found", "core memory not set for this agent"))
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]any{
|
||||
"owner_id": ownerStr,
|
||||
"agent_name": agentName,
|
||||
"blob": blob,
|
||||
"updated_at": updatedAt.Format("2006-01-02T15:04:05Z07:00"),
|
||||
})
|
||||
}
|
||||
|
||||
// Put handles PUT /api/owner/{ownerID}/agents/{agentName}/core-memory.
|
||||
// Body: {"blob": "..."}. Returns 200 on success, 413 on
|
||||
// core_memory_too_large, 400 on malformed body.
|
||||
func (h *MemoryCoreHandler) Put(w http.ResponseWriter, r *http.Request) {
|
||||
if h.store == nil {
|
||||
writeJSON(w, http.StatusServiceUnavailable, errorBody("unavailable", "core memory store not configured"))
|
||||
return
|
||||
}
|
||||
ownerStr, agentName, ok := h.authorize(w, r)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
var body struct {
|
||||
Blob string `json:"blob"`
|
||||
UpdatedBy string `json:"updated_by,omitempty"`
|
||||
}
|
||||
if err := json.NewDecoder(r.Body).Decode(&body); err != nil {
|
||||
writeJSON(w, http.StatusBadRequest, errorBody("bad_request", "Invalid JSON body"))
|
||||
return
|
||||
}
|
||||
updatedBy := body.UpdatedBy
|
||||
if updatedBy == "" {
|
||||
updatedBy = "human"
|
||||
}
|
||||
if err := h.store.Set(r.Context(), ownerStr, agentName, body.Blob, updatedBy); err != nil {
|
||||
if errors.Is(err, messaging.ErrCoreMemoryTooLarge) {
|
||||
writeJSON(w, http.StatusRequestEntityTooLarge, errorBody(
|
||||
"core_memory_too_large",
|
||||
fmt.Sprintf("Blob exceeds %d bytes", h.store.MaxBytes()),
|
||||
))
|
||||
return
|
||||
}
|
||||
h.logger.Error("memory_core set failed", "error", err, "owner", ownerStr, "agent", agentName)
|
||||
writeJSON(w, http.StatusInternalServerError, errorBody("server_error", "Failed to write core memory"))
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]any{
|
||||
"owner_id": ownerStr,
|
||||
"agent_name": agentName,
|
||||
"blob_chars": len(body.Blob),
|
||||
"updated_by": updatedBy,
|
||||
})
|
||||
}
|
||||
|
||||
// Delete handles DELETE /api/owner/{ownerID}/agents/{agentName}/core-memory.
|
||||
// Returns 204 when a row existed and was removed, 404 when no row was
|
||||
// present at the start of the call.
|
||||
func (h *MemoryCoreHandler) Delete(w http.ResponseWriter, r *http.Request) {
|
||||
if h.store == nil {
|
||||
writeJSON(w, http.StatusServiceUnavailable, errorBody("unavailable", "core memory store not configured"))
|
||||
return
|
||||
}
|
||||
ownerStr, agentName, ok := h.authorize(w, r)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
// Check existence so we can return the canonical 204 vs 404. The
|
||||
// store's Delete is idempotent — it never errors on missing rows.
|
||||
_, _, exists, err := h.store.Get(r.Context(), ownerStr, agentName)
|
||||
if err != nil {
|
||||
h.logger.Error("memory_core get-before-delete failed", "error", err, "owner", ownerStr, "agent", agentName)
|
||||
writeJSON(w, http.StatusInternalServerError, errorBody("server_error", "Failed to read core memory"))
|
||||
return
|
||||
}
|
||||
if !exists {
|
||||
writeJSON(w, http.StatusNotFound, errorBody("not_found", "core memory not set for this agent"))
|
||||
return
|
||||
}
|
||||
if err := h.store.Delete(r.Context(), ownerStr, agentName); err != nil {
|
||||
h.logger.Error("memory_core delete failed", "error", err, "owner", ownerStr, "agent", agentName)
|
||||
writeJSON(w, http.StatusInternalServerError, errorBody("server_error", "Failed to delete core memory"))
|
||||
return
|
||||
}
|
||||
w.WriteHeader(http.StatusNoContent)
|
||||
}
|
||||
+17
-1
@@ -15,9 +15,9 @@ import (
|
||||
"github.com/synapbus/synapbus/internal/harness/runs"
|
||||
"github.com/synapbus/synapbus/internal/k8s"
|
||||
"github.com/synapbus/synapbus/internal/messaging"
|
||||
"github.com/synapbus/synapbus/internal/reactor"
|
||||
"github.com/synapbus/synapbus/internal/push"
|
||||
"github.com/synapbus/synapbus/internal/reactions"
|
||||
"github.com/synapbus/synapbus/internal/reactor"
|
||||
"github.com/synapbus/synapbus/internal/trace"
|
||||
"github.com/synapbus/synapbus/internal/trust"
|
||||
"github.com/synapbus/synapbus/internal/webhooks"
|
||||
@@ -54,6 +54,10 @@ type RouterConfig struct {
|
||||
DB *sql.DB
|
||||
Version string
|
||||
BaseURL string
|
||||
|
||||
// CoreMemoryStore (feature 020 — US2) wires the per-(owner, agent)
|
||||
// core memory REST endpoints. Nil → routes not registered.
|
||||
CoreMemoryStore *messaging.CoreMemoryStore
|
||||
}
|
||||
|
||||
// NewRouter creates a chi router with all API routes configured.
|
||||
@@ -333,6 +337,18 @@ func NewRouterWithConfig(cfg RouterConfig) chi.Router {
|
||||
})
|
||||
}
|
||||
|
||||
// Per-(owner, agent) core memory (feature 020 — US2)
|
||||
if cfg.CoreMemoryStore != nil {
|
||||
coreHandler := NewMemoryCoreHandler(cfg.CoreMemoryStore)
|
||||
r.Group(func(r chi.Router) {
|
||||
r.Use(authMiddleware)
|
||||
|
||||
r.Get("/api/owner/{ownerID}/agents/{agentName}/core-memory", coreHandler.Get)
|
||||
r.Put("/api/owner/{ownerID}/agents/{agentName}/core-memory", coreHandler.Put)
|
||||
r.Delete("/api/owner/{ownerID}/agents/{agentName}/core-memory", coreHandler.Delete)
|
||||
})
|
||||
}
|
||||
|
||||
// Version (unauthenticated)
|
||||
if cfg.Version != "" {
|
||||
versionHandler := NewVersionHandler(cfg.Version)
|
||||
|
||||
@@ -0,0 +1,158 @@
|
||||
// Integration test for US2 → US1 wiring: a real CoreMemoryStore, wired
|
||||
// through messaging.NewCoreProvider, surfaces a seeded blob in the
|
||||
// wrapped tool response's `relevant_context.core_memory` field; an agent
|
||||
// with no row gets no relevant_context.
|
||||
package mcp
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"testing"
|
||||
|
||||
mcplib "github.com/mark3labs/mcp-go/mcp"
|
||||
_ "modernc.org/sqlite"
|
||||
|
||||
"github.com/synapbus/synapbus/internal/agents"
|
||||
"github.com/synapbus/synapbus/internal/messaging"
|
||||
"github.com/synapbus/synapbus/internal/storage"
|
||||
)
|
||||
|
||||
func newInjectionTestDB(t *testing.T) *sql.DB {
|
||||
t.Helper()
|
||||
dsn := fmt.Sprintf("file:%s?mode=memory&cache=shared", t.Name())
|
||||
db, err := sql.Open("sqlite", dsn)
|
||||
if err != nil {
|
||||
t.Fatalf("open db: %v", err)
|
||||
}
|
||||
t.Cleanup(func() { db.Close() })
|
||||
if _, err := db.Exec("PRAGMA foreign_keys=ON"); err != nil {
|
||||
t.Fatalf("foreign keys: %v", err)
|
||||
}
|
||||
if err := storage.RunMigrations(context.Background(), db); err != nil {
|
||||
t.Fatalf("migrations: %v", err)
|
||||
}
|
||||
return db
|
||||
}
|
||||
|
||||
// TestInjection_CoreMemoryWiring proves a wrapped session-start handler
|
||||
// surfaces a blob seeded via CoreMemoryStore as
|
||||
// `relevant_context.core_memory`. Mirrors the contract example in
|
||||
// `specs/020-proactive-memory-dream-worker/contracts/mcp-injection.md`.
|
||||
func TestInjection_CoreMemoryWiring(t *testing.T) {
|
||||
db := newInjectionTestDB(t)
|
||||
ctx := context.Background()
|
||||
|
||||
const owner = "42"
|
||||
const agent = "a1"
|
||||
const blob = "I am a1. Currently focused on memory tests."
|
||||
|
||||
coreStore := messaging.NewCoreMemoryStore(db, 2048)
|
||||
if err := coreStore.Set(ctx, owner, agent, blob, "human"); err != nil {
|
||||
t.Fatalf("seed core memory: %v", err)
|
||||
}
|
||||
|
||||
cfg := WrapConfig{
|
||||
Cfg: messaging.MemoryConfig{
|
||||
InjectionEnabled: true,
|
||||
InjectionBudgetTokens: 500,
|
||||
InjectionMaxItems: 5,
|
||||
InjectionMinScore: 0.25,
|
||||
},
|
||||
SearchSvc: nil, // query="" forces no retrieval — only core matters.
|
||||
IncludeCore: true,
|
||||
CoreProvider: messaging.NewCoreProvider(coreStore),
|
||||
QuerySource: func(_ context.Context, _ string, _ map[string]any, _ map[string]any) string { return "" },
|
||||
}
|
||||
inner := stubHandler(map[string]any{"agent": agent})
|
||||
wrapped := WrapInjection(inner, "my_status", cfg)
|
||||
|
||||
// Owner 42 ↔ caller a1.
|
||||
callerCtx := agents.ContextWithAgent(ctx, &agents.Agent{Name: agent, OwnerID: 42})
|
||||
res, err := wrapped(callerCtx, mcplib.CallToolRequest{})
|
||||
if err != nil {
|
||||
t.Fatalf("wrapped my_status: %v", err)
|
||||
}
|
||||
body := extractJSON(t, res)
|
||||
rc, ok := body["relevant_context"].(map[string]any)
|
||||
if !ok {
|
||||
t.Fatalf("relevant_context missing: %+v", body)
|
||||
}
|
||||
if got := rc["core_memory"]; got != blob {
|
||||
t.Errorf("core_memory: got %v want %q", got, blob)
|
||||
}
|
||||
}
|
||||
|
||||
// TestInjection_NoCoreMemoryYieldsNoPacket verifies that an agent
|
||||
// without a memory_core row gets the original handler response back,
|
||||
// without a `relevant_context` field appended.
|
||||
func TestInjection_NoCoreMemoryYieldsNoPacket(t *testing.T) {
|
||||
db := newInjectionTestDB(t)
|
||||
|
||||
coreStore := messaging.NewCoreMemoryStore(db, 2048)
|
||||
// Intentionally NO Set — the agent has no row.
|
||||
|
||||
cfg := WrapConfig{
|
||||
Cfg: messaging.MemoryConfig{
|
||||
InjectionEnabled: true,
|
||||
InjectionBudgetTokens: 500,
|
||||
InjectionMaxItems: 5,
|
||||
InjectionMinScore: 0.25,
|
||||
},
|
||||
SearchSvc: nil,
|
||||
IncludeCore: true,
|
||||
CoreProvider: messaging.NewCoreProvider(coreStore),
|
||||
QuerySource: func(_ context.Context, _ string, _ map[string]any, _ map[string]any) string { return "" },
|
||||
}
|
||||
inner := stubHandler(map[string]any{"agent": "a2"})
|
||||
wrapped := WrapInjection(inner, "my_status", cfg)
|
||||
|
||||
ctx := agents.ContextWithAgent(context.Background(), &agents.Agent{Name: "a2", OwnerID: 42})
|
||||
res, err := wrapped(ctx, mcplib.CallToolRequest{})
|
||||
if err != nil {
|
||||
t.Fatalf("wrapped: %v", err)
|
||||
}
|
||||
body := extractJSON(t, res)
|
||||
if _, has := body["relevant_context"]; has {
|
||||
t.Errorf("relevant_context attached for agent with no core row: %v", body["relevant_context"])
|
||||
}
|
||||
if body["agent"] != "a2" {
|
||||
t.Errorf("inner body lost: %+v", body)
|
||||
}
|
||||
}
|
||||
|
||||
// Verify ContextPacket round-trips its core_memory through json. This is
|
||||
// the contract field consumed by clients.
|
||||
func TestInjection_CoreMemoryJSONShape(t *testing.T) {
|
||||
db := newInjectionTestDB(t)
|
||||
ctx := context.Background()
|
||||
coreStore := messaging.NewCoreMemoryStore(db, 2048)
|
||||
if err := coreStore.Set(ctx, "1", "a1", "hello", "human"); err != nil {
|
||||
t.Fatalf("seed: %v", err)
|
||||
}
|
||||
|
||||
provider := messaging.NewCoreProvider(coreStore)
|
||||
got, err := provider.Get(ctx, "1", "a1")
|
||||
if err != nil {
|
||||
t.Fatalf("provider.Get: %v", err)
|
||||
}
|
||||
if got != "hello" {
|
||||
t.Errorf("provider.Get: got %q want hello", got)
|
||||
}
|
||||
|
||||
// Empty case (no row) yields "" without error.
|
||||
got, err = provider.Get(ctx, "1", "nobody")
|
||||
if err != nil || got != "" {
|
||||
t.Errorf("provider.Get on missing: got %q err=%v", got, err)
|
||||
}
|
||||
|
||||
// Sanity: ensure the adapter is reachable through json marshaling of a packet.
|
||||
type fakePacket struct {
|
||||
Core string `json:"core_memory,omitempty"`
|
||||
}
|
||||
b, _ := json.Marshal(fakePacket{Core: "hello"})
|
||||
if string(b) != `{"core_memory":"hello"}` {
|
||||
t.Errorf("json marshaling: got %s", b)
|
||||
}
|
||||
}
|
||||
+26
-7
@@ -29,13 +29,13 @@ import (
|
||||
|
||||
// MCPServer wraps the mcp-go server with SynapBus services.
|
||||
type MCPServer struct {
|
||||
mcpServer *server.MCPServer
|
||||
httpServer *server.StreamableHTTPServer
|
||||
connMgr *ConnectionManager
|
||||
agentService *agents.AgentService
|
||||
hybridRegistrar *HybridToolRegistrar
|
||||
logger *slog.Logger
|
||||
console *console.Printer
|
||||
mcpServer *server.MCPServer
|
||||
httpServer *server.StreamableHTTPServer
|
||||
connMgr *ConnectionManager
|
||||
agentService *agents.AgentService
|
||||
hybridRegistrar *HybridToolRegistrar
|
||||
logger *slog.Logger
|
||||
console *console.Printer
|
||||
}
|
||||
|
||||
// NewMCPServer creates and configures a new MCP server with 5 hybrid tools registered.
|
||||
@@ -206,6 +206,25 @@ func NewMCPServer(
|
||||
return s
|
||||
}
|
||||
|
||||
// SetInjection wires the proactive-memory injection middleware
|
||||
// (feature 020) into the hybrid tool registrar and re-registers the
|
||||
// hybrid tools so the wrappers take effect. Must be called after
|
||||
// NewMCPServer and before the server starts handling traffic.
|
||||
//
|
||||
// `coreProvider` is consulted only on session-start tools (currently
|
||||
// `my_status`). Pass nil when US2 has not yet been wired.
|
||||
func (s *MCPServer) SetInjection(cfg messaging.MemoryConfig, store *messaging.MemoryInjections, coreProvider search.CoreMemoryProvider) {
|
||||
if s.hybridRegistrar == nil || s.mcpServer == nil {
|
||||
return
|
||||
}
|
||||
s.hybridRegistrar.SetInjection(cfg, store, coreProvider)
|
||||
// Re-register the hybrid tools so the new InjectionEnabled / Core
|
||||
// wiring takes effect. AddTool overwrites by name (see mcp-go's
|
||||
// `MCPServer.AddTools`), so this swaps in the wrapped handlers
|
||||
// without leaking the original registrations.
|
||||
s.hybridRegistrar.RegisterAllOnServer(s.mcpServer)
|
||||
}
|
||||
|
||||
// WireGoalsTools registers the spec-018 tool surface (create_goal,
|
||||
// propose_task_tree, claim_task, request_resource, list_resources,
|
||||
// complete_goal) on the MCP server. Must be called after NewMCPServer.
|
||||
|
||||
@@ -0,0 +1,195 @@
|
||||
// Per-(owner, agent) core memory store. Backs User Story 2 of feature
|
||||
// 020-proactive-memory-dream-worker — small, owner-scoped, replace-wholesale
|
||||
// blobs surfaced in `relevant_context.core_memory` on session-start tools
|
||||
// (e.g. `my_status`).
|
||||
//
|
||||
// Schema lives in `internal/storage/schema/028_memory_consolidation.sql`
|
||||
// (table `memory_core`). owner_id is stored as TEXT (string form of
|
||||
// `users.id`) to match the proactive-memory tables and the request-context
|
||||
// owner_id propagated by auth middleware.
|
||||
package messaging
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"errors"
|
||||
"fmt"
|
||||
"time"
|
||||
)
|
||||
|
||||
// ErrCoreMemoryTooLarge is returned by CoreMemoryStore.Set when the blob
|
||||
// exceeds the configured max bytes (default SYNAPBUS_CORE_MEMORY_MAX_BYTES = 2048).
|
||||
// The MCP `memory_rewrite_core` tool surfaces this as the contractual
|
||||
// `core_memory_too_large` error code (see contracts/mcp-memory-tools.md).
|
||||
var ErrCoreMemoryTooLarge = errors.New("core memory blob exceeds max bytes")
|
||||
|
||||
// CoreMemoryRecord is one row of the `memory_core` table.
|
||||
type CoreMemoryRecord struct {
|
||||
OwnerID string
|
||||
AgentName string
|
||||
Blob string
|
||||
UpdatedAt time.Time
|
||||
UpdatedBy string
|
||||
}
|
||||
|
||||
// CoreMemoryStore wraps the `memory_core` table. All operations are
|
||||
// owner-scoped — owner_id is part of the primary key — so callers cannot
|
||||
// cross-read another owner's blobs.
|
||||
type CoreMemoryStore struct {
|
||||
db *sql.DB
|
||||
maxBytes int
|
||||
}
|
||||
|
||||
// NewCoreMemoryStore returns a store rooted at db enforcing the given
|
||||
// max-bytes cap on Set. When maxBytes <= 0, defaults to 2048 (the spec
|
||||
// default for SYNAPBUS_CORE_MEMORY_MAX_BYTES).
|
||||
func NewCoreMemoryStore(db *sql.DB, maxBytes int) *CoreMemoryStore {
|
||||
if maxBytes <= 0 {
|
||||
maxBytes = 2048
|
||||
}
|
||||
return &CoreMemoryStore{db: db, maxBytes: maxBytes}
|
||||
}
|
||||
|
||||
// MaxBytes returns the configured upper bound for Set blobs.
|
||||
func (s *CoreMemoryStore) MaxBytes() int { return s.maxBytes }
|
||||
|
||||
// Get returns the core memory blob for (ownerID, agentName). When no row
|
||||
// exists, returns ok=false with no error. Unexpected DB errors surface as
|
||||
// err.
|
||||
func (s *CoreMemoryStore) Get(ctx context.Context, ownerID, agentName string) (blob string, updatedAt time.Time, ok bool, err error) {
|
||||
row := s.db.QueryRowContext(ctx,
|
||||
`SELECT blob, updated_at FROM memory_core WHERE owner_id = ? AND agent_name = ?`,
|
||||
ownerID, agentName,
|
||||
)
|
||||
if err := row.Scan(&blob, &updatedAt); err != nil {
|
||||
if errors.Is(err, sql.ErrNoRows) {
|
||||
return "", time.Time{}, false, nil
|
||||
}
|
||||
return "", time.Time{}, false, fmt.Errorf("core memory get: %w", err)
|
||||
}
|
||||
return blob, updatedAt, true, nil
|
||||
}
|
||||
|
||||
// Set wholesale-replaces the blob for (ownerID, agentName). Enforces the
|
||||
// configured size cap and returns ErrCoreMemoryTooLarge when violated.
|
||||
// `updatedBy` is recorded for audit (typically the caller agent name, or
|
||||
// "human" for admin/web edits).
|
||||
func (s *CoreMemoryStore) Set(ctx context.Context, ownerID, agentName, blob, updatedBy string) error {
|
||||
if len(blob) > s.maxBytes {
|
||||
return ErrCoreMemoryTooLarge
|
||||
}
|
||||
if ownerID == "" {
|
||||
return fmt.Errorf("core memory set: empty owner_id")
|
||||
}
|
||||
if agentName == "" {
|
||||
return fmt.Errorf("core memory set: empty agent_name")
|
||||
}
|
||||
if updatedBy == "" {
|
||||
return fmt.Errorf("core memory set: empty updated_by")
|
||||
}
|
||||
_, err := s.db.ExecContext(ctx,
|
||||
`INSERT INTO memory_core (owner_id, agent_name, blob, updated_at, updated_by)
|
||||
VALUES (?, ?, ?, CURRENT_TIMESTAMP, ?)
|
||||
ON CONFLICT(owner_id, agent_name) DO UPDATE SET
|
||||
blob = excluded.blob,
|
||||
updated_at = CURRENT_TIMESTAMP,
|
||||
updated_by = excluded.updated_by`,
|
||||
ownerID, agentName, blob, updatedBy,
|
||||
)
|
||||
if err != nil {
|
||||
return fmt.Errorf("core memory set: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Delete removes the (ownerID, agentName) row. Returns nil even if no
|
||||
// row matched — callers should treat "deleted" and "did not exist" the
|
||||
// same way (the REST endpoint distinguishes via a separate Get).
|
||||
func (s *CoreMemoryStore) Delete(ctx context.Context, ownerID, agentName string) error {
|
||||
_, err := s.db.ExecContext(ctx,
|
||||
`DELETE FROM memory_core WHERE owner_id = ? AND agent_name = ?`,
|
||||
ownerID, agentName,
|
||||
)
|
||||
if err != nil {
|
||||
return fmt.Errorf("core memory delete: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// List returns all core memory rows for the given owner. Used by the
|
||||
// future audit UI (deferred US4) and by admin tooling.
|
||||
func (s *CoreMemoryStore) List(ctx context.Context, ownerID string) ([]CoreMemoryRecord, error) {
|
||||
rows, err := s.db.QueryContext(ctx,
|
||||
`SELECT owner_id, agent_name, blob, updated_at, updated_by
|
||||
FROM memory_core
|
||||
WHERE owner_id = ?
|
||||
ORDER BY agent_name`,
|
||||
ownerID,
|
||||
)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("core memory list: %w", err)
|
||||
}
|
||||
defer rows.Close()
|
||||
var out []CoreMemoryRecord
|
||||
for rows.Next() {
|
||||
var r CoreMemoryRecord
|
||||
if err := rows.Scan(&r.OwnerID, &r.AgentName, &r.Blob, &r.UpdatedAt, &r.UpdatedBy); err != nil {
|
||||
return nil, fmt.Errorf("core memory list scan: %w", err)
|
||||
}
|
||||
out = append(out, r)
|
||||
}
|
||||
if err := rows.Err(); err != nil {
|
||||
return nil, fmt.Errorf("core memory list rows: %w", err)
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// GetForInjection is the adapter implementing
|
||||
// `search.CoreMemoryProvider.Get`. Returns "" (no error) when no row
|
||||
// exists — the empty-string convention lets the injection wrapper treat a
|
||||
// missing core memory the same as "field omitted", per
|
||||
// `contracts/mcp-injection.md`.
|
||||
//
|
||||
// The matching interface contract is in
|
||||
// `internal/search/injection.go`'s `CoreMemoryProvider`:
|
||||
//
|
||||
// Get(ctx context.Context, ownerID, agentName string) (string, error)
|
||||
//
|
||||
// We expose this as a method on the store (not a separate type) so
|
||||
// callers can pass `coreStore.GetForInjection` as a method value — but the
|
||||
// store itself also satisfies the interface via its `Get` method below.
|
||||
func (s *CoreMemoryStore) GetForInjection(ctx context.Context, ownerID, agentName string) (string, error) {
|
||||
blob, _, ok, err := s.Get(ctx, ownerID, agentName)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
if !ok {
|
||||
return "", nil
|
||||
}
|
||||
return blob, nil
|
||||
}
|
||||
|
||||
// coreProviderAdapter wraps a *CoreMemoryStore so it satisfies
|
||||
// `search.CoreMemoryProvider` (which requires a `Get(ctx, ownerID,
|
||||
// agentName) (string, error)` signature — distinct from the store's
|
||||
// 4-return Get). Use NewCoreProvider to construct.
|
||||
type coreProviderAdapter struct {
|
||||
store *CoreMemoryStore
|
||||
}
|
||||
|
||||
// NewCoreProvider returns an object satisfying
|
||||
// `search.CoreMemoryProvider` so callers can wire the store into
|
||||
// WrapConfig.CoreProvider without leaking the store's richer Get
|
||||
// signature.
|
||||
func NewCoreProvider(store *CoreMemoryStore) *coreProviderAdapter {
|
||||
return &coreProviderAdapter{store: store}
|
||||
}
|
||||
|
||||
// Get implements `search.CoreMemoryProvider.Get`. Returns "" when no row
|
||||
// exists for the (owner, agent) pair.
|
||||
func (a *coreProviderAdapter) Get(ctx context.Context, ownerID, agentName string) (string, error) {
|
||||
if a == nil || a.store == nil {
|
||||
return "", nil
|
||||
}
|
||||
return a.store.GetForInjection(ctx, ownerID, agentName)
|
||||
}
|
||||
@@ -0,0 +1,179 @@
|
||||
package messaging
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestCoreMemoryStore_GetSetRoundTrip(t *testing.T) {
|
||||
db := newTestDB(t)
|
||||
store := NewCoreMemoryStore(db, 2048)
|
||||
ctx := context.Background()
|
||||
|
||||
const owner = "1"
|
||||
const agent = "research-mcpproxy"
|
||||
const blob = "You are research-mcpproxy. Focus on benchmarking."
|
||||
|
||||
if err := store.Set(ctx, owner, agent, blob, "human"); err != nil {
|
||||
t.Fatalf("Set: %v", err)
|
||||
}
|
||||
|
||||
got, updatedAt, ok, err := store.Get(ctx, owner, agent)
|
||||
if err != nil {
|
||||
t.Fatalf("Get: %v", err)
|
||||
}
|
||||
if !ok {
|
||||
t.Fatal("Get: expected ok=true, got false")
|
||||
}
|
||||
if got != blob {
|
||||
t.Errorf("Get blob mismatch: got %q want %q", got, blob)
|
||||
}
|
||||
if updatedAt.IsZero() {
|
||||
t.Error("Get: expected non-zero updated_at")
|
||||
}
|
||||
}
|
||||
|
||||
func TestCoreMemoryStore_SetReplacesWholesale(t *testing.T) {
|
||||
db := newTestDB(t)
|
||||
store := NewCoreMemoryStore(db, 2048)
|
||||
ctx := context.Background()
|
||||
|
||||
const owner = "1"
|
||||
const agent = "alpha"
|
||||
|
||||
if err := store.Set(ctx, owner, agent, "first version", "human"); err != nil {
|
||||
t.Fatalf("first Set: %v", err)
|
||||
}
|
||||
if err := store.Set(ctx, owner, agent, "second version", "dream:1"); err != nil {
|
||||
t.Fatalf("second Set: %v", err)
|
||||
}
|
||||
|
||||
got, _, ok, err := store.Get(ctx, owner, agent)
|
||||
if err != nil {
|
||||
t.Fatalf("Get: %v", err)
|
||||
}
|
||||
if !ok || got != "second version" {
|
||||
t.Errorf("expected wholesale replace, got ok=%v blob=%q", ok, got)
|
||||
}
|
||||
if strings.Contains(got, "first") {
|
||||
t.Errorf("expected merge to NOT happen; got %q", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCoreMemoryStore_SetRejectsOverSize(t *testing.T) {
|
||||
db := newTestDB(t)
|
||||
store := NewCoreMemoryStore(db, 16)
|
||||
ctx := context.Background()
|
||||
|
||||
tooBig := strings.Repeat("x", 17)
|
||||
err := store.Set(ctx, "1", "alpha", tooBig, "human")
|
||||
if err == nil {
|
||||
t.Fatal("Set: expected ErrCoreMemoryTooLarge, got nil")
|
||||
}
|
||||
if !errors.Is(err, ErrCoreMemoryTooLarge) {
|
||||
t.Errorf("Set: expected ErrCoreMemoryTooLarge, got %v", err)
|
||||
}
|
||||
|
||||
// Exactly at-cap is fine.
|
||||
atCap := strings.Repeat("x", 16)
|
||||
if err := store.Set(ctx, "1", "alpha", atCap, "human"); err != nil {
|
||||
t.Errorf("Set at cap: unexpected error %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCoreMemoryStore_Delete(t *testing.T) {
|
||||
db := newTestDB(t)
|
||||
store := NewCoreMemoryStore(db, 2048)
|
||||
ctx := context.Background()
|
||||
|
||||
if err := store.Set(ctx, "1", "alpha", "blob", "human"); err != nil {
|
||||
t.Fatalf("Set: %v", err)
|
||||
}
|
||||
if err := store.Delete(ctx, "1", "alpha"); err != nil {
|
||||
t.Fatalf("Delete: %v", err)
|
||||
}
|
||||
_, _, ok, err := store.Get(ctx, "1", "alpha")
|
||||
if err != nil {
|
||||
t.Fatalf("Get after Delete: %v", err)
|
||||
}
|
||||
if ok {
|
||||
t.Error("Get after Delete: expected ok=false")
|
||||
}
|
||||
|
||||
// Delete on missing row is idempotent.
|
||||
if err := store.Delete(ctx, "1", "alpha"); err != nil {
|
||||
t.Errorf("Delete on missing: unexpected error %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCoreMemoryStore_OwnerScoping(t *testing.T) {
|
||||
db := newTestDB(t)
|
||||
store := NewCoreMemoryStore(db, 2048)
|
||||
ctx := context.Background()
|
||||
|
||||
if err := store.Set(ctx, "1", "shared-name", "H1 blob", "h1"); err != nil {
|
||||
t.Fatalf("Set H1: %v", err)
|
||||
}
|
||||
if err := store.Set(ctx, "2", "shared-name", "H2 blob", "h2"); err != nil {
|
||||
t.Fatalf("Set H2: %v", err)
|
||||
}
|
||||
|
||||
got1, _, ok1, err := store.Get(ctx, "1", "shared-name")
|
||||
if err != nil || !ok1 || got1 != "H1 blob" {
|
||||
t.Errorf("H1 Get: ok=%v got=%q err=%v", ok1, got1, err)
|
||||
}
|
||||
got2, _, ok2, err := store.Get(ctx, "2", "shared-name")
|
||||
if err != nil || !ok2 || got2 != "H2 blob" {
|
||||
t.Errorf("H2 Get: ok=%v got=%q err=%v", ok2, got2, err)
|
||||
}
|
||||
|
||||
// H1 cannot see H2's blob and vice versa — distinct PKs.
|
||||
if got1 == got2 {
|
||||
t.Error("owner scoping broken: H1 and H2 see the same blob")
|
||||
}
|
||||
|
||||
// List is scoped to a single owner.
|
||||
listH1, err := store.List(ctx, "1")
|
||||
if err != nil {
|
||||
t.Fatalf("List H1: %v", err)
|
||||
}
|
||||
if len(listH1) != 1 || listH1[0].OwnerID != "1" {
|
||||
t.Errorf("List H1: expected 1 row owned by 1, got %#v", listH1)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCoreMemoryStore_GetForInjection(t *testing.T) {
|
||||
db := newTestDB(t)
|
||||
store := NewCoreMemoryStore(db, 2048)
|
||||
ctx := context.Background()
|
||||
|
||||
// Missing row → "" + no error (NOT sql.ErrNoRows).
|
||||
blob, err := store.GetForInjection(ctx, "1", "missing-agent")
|
||||
if err != nil {
|
||||
t.Fatalf("GetForInjection on missing: unexpected error %v", err)
|
||||
}
|
||||
if blob != "" {
|
||||
t.Errorf("GetForInjection on missing: expected \"\" got %q", blob)
|
||||
}
|
||||
|
||||
// Present row → blob.
|
||||
if err := store.Set(ctx, "1", "alpha", "core blob", "human"); err != nil {
|
||||
t.Fatalf("Set: %v", err)
|
||||
}
|
||||
blob, err = store.GetForInjection(ctx, "1", "alpha")
|
||||
if err != nil {
|
||||
t.Fatalf("GetForInjection: %v", err)
|
||||
}
|
||||
if blob != "core blob" {
|
||||
t.Errorf("GetForInjection: got %q want %q", blob, "core blob")
|
||||
}
|
||||
|
||||
// Adapter satisfies search.CoreMemoryProvider implicitly.
|
||||
provider := NewCoreProvider(store)
|
||||
got, err := provider.Get(ctx, "1", "alpha")
|
||||
if err != nil || got != "core blob" {
|
||||
t.Errorf("CoreProvider.Get: got %q err=%v", got, err)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user