Drain-on-demand: SYNAPBUS_DREAM_PARALLEL (default 1) and `synapbus memory dream-run --parallel N` fan out N concurrent dream-agent k8s Jobs per (owner, job_type) in one shot. Set high (e.g. 8) to drain backlog quickly, then back to 1 for normal hourly operation. Schema: - migration 030_dream_parallelism: adds slot INTEGER NOT NULL DEFAULT 0 to memory_consolidation_jobs. Drops + recreates the partial unique in-flight index as (owner, job_type, slot) so slots 0..N-1 each hold one in-flight job independently. Stores: - JobsStore.CreateOnSlot + CreateNextAvailableSlot. - ConsolidatorWorker.ForceRunN dispatches N parallel jobs through the existing launchOne path (extracted from ForceRun). - core_rewrite coerces to N=1 regardless of the knob — per-(owner, agent) blob is wholesale-replace and concurrent rewrites would race. Three bug fixes discovered while bringing the parallel path up on kubic: 1. k8s Job names collided on rapid relaunch because runner.go used "synapbus-<agent>-<msg_id>", and dream dispatches have msg_id=0. Now appends a unique (timestamp%1e6, 4-byte random) suffix when msg_id is zero; historical "synapbus-<agent>-<id>" prefix preserved. 2. memory_list_unprocessed didn't actually exclude already-refined messages — the contract said it should, the implementation returned the same oldest-50 every cycle. The dream agent kept re-refining the same set: 221 refines links touched only 55 unique dst messages, so progress flat-lined. Added the NOT IN (refines/duplicate_of/superseded_by) filter and a from_agent NOT LIKE 'dream:%' clause so the agent never refines its own reflections. 3. The k8sjob harness was constructed with nil Waiter in main.go, so every dream dispatch failed instantly with "k8sjob: no Waiter configured". Now builds a ClientsetWaiter from the in-cluster clientset. Plus admin/server.go gets DreamRunN closure + DefaultDreamParallel (sourced from MemoryConfig.DreamParallel). admin/socket.go handleMemoryDreamRun accepts `parallel` arg and returns job_ids[]. CLI admin command grows --parallel N flag. Live evidence from kubic (image v0.21.0-amd64): 1 CLI call with --parallel 8 produced 8 job rows on slots 0..7, spawned 8 distinct k8s Jobs with unique suffixes, retired ~86 unprocessed messages in <1 min (vs ~10/cycle for the buggy serial version pre-fix-2). Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
104 lines
3.9 KiB
Go
104 lines
3.9 KiB
Go
// Package admin provides a Unix domain socket server for local administration.
|
|
package admin
|
|
|
|
import (
|
|
"context"
|
|
"database/sql"
|
|
"log/slog"
|
|
"net"
|
|
|
|
"github.com/synapbus/synapbus/internal/agents"
|
|
"github.com/synapbus/synapbus/internal/attachments"
|
|
"github.com/synapbus/synapbus/internal/auth"
|
|
"github.com/synapbus/synapbus/internal/channels"
|
|
"github.com/synapbus/synapbus/internal/k8s"
|
|
"github.com/synapbus/synapbus/internal/messaging"
|
|
"github.com/synapbus/synapbus/internal/search"
|
|
"github.com/synapbus/synapbus/internal/trace"
|
|
"github.com/synapbus/synapbus/internal/webhooks"
|
|
)
|
|
|
|
// WebhookServiceProvider defines the webhook operations needed by the admin socket.
|
|
type WebhookServiceProvider interface {
|
|
RegisterWebhook(ctx context.Context, agentName, url string, events []string, secret string) (*webhooks.Webhook, error)
|
|
ListWebhooks(ctx context.Context, agentName string) ([]*webhooks.Webhook, error)
|
|
DeleteWebhook(ctx context.Context, agentName string, webhookID int64) error
|
|
}
|
|
|
|
// K8sServiceProvider defines the K8s handler operations needed by the admin socket.
|
|
type K8sServiceProvider interface {
|
|
RegisterHandler(ctx context.Context, agentName string, req k8s.RegisterHandlerRequest) (*k8s.K8sHandler, error)
|
|
ListHandlers(ctx context.Context, agentName string) ([]*k8s.K8sHandler, error)
|
|
DeleteHandler(ctx context.Context, agentName string, handlerID int64) error
|
|
}
|
|
|
|
// 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
|
|
|
|
// 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
|
|
|
|
// DreamRun, when non-nil, dispatches a single consolidation job
|
|
// bypassing the trigger check. Wired by main.go when the
|
|
// consolidator worker is enabled. Closure form keeps the worker
|
|
// internals out of the admin package's import graph.
|
|
DreamRun func(ctx context.Context, ownerID, jobType string) (jobID int64, err error)
|
|
|
|
// DreamRunN fans out N parallel consolidation jobs (via slot 0..N-1)
|
|
// for one (owner, job_type). Used by `synapbus memory dream-run
|
|
// --parallel N`. core_rewrite always coerces to N=1 server-side.
|
|
DreamRunN func(ctx context.Context, ownerID, jobType string, parallel int) (jobIDs []int64, err error)
|
|
|
|
// DefaultDreamParallel is consulted when the CLI request omits
|
|
// --parallel. Sourced from MemoryConfig.DreamParallel.
|
|
DefaultDreamParallel int
|
|
}
|
|
|
|
// RetentionStatusProvider provides retention status information.
|
|
type RetentionStatusProvider interface {
|
|
Status() map[string]interface{}
|
|
}
|
|
|
|
// AdminServer is a Unix domain socket server for local administration.
|
|
type AdminServer struct {
|
|
listener net.Listener
|
|
db *sql.DB
|
|
services *Services
|
|
logger *slog.Logger
|
|
socketPath string
|
|
done chan struct{}
|
|
}
|
|
|
|
// NewServer creates a new admin server bound to a Unix socket at {dataDir}/synapbus.sock.
|
|
// If socketPath is non-empty it overrides the default.
|
|
func NewServer(socketPath string, db *sql.DB, services *Services, logger *slog.Logger) *AdminServer {
|
|
return &AdminServer{
|
|
db: db,
|
|
services: services,
|
|
logger: logger.With("component", "admin"),
|
|
socketPath: socketPath,
|
|
done: make(chan struct{}),
|
|
}
|
|
}
|
|
|
|
// SocketPath returns the path to the Unix socket.
|
|
func (s *AdminServer) SocketPath() string {
|
|
return s.socketPath
|
|
}
|