diff --git a/cmd/synapbus/admin.go b/cmd/synapbus/admin.go index 82d0079..272d296 100644 --- a/cmd/synapbus/admin.go +++ b/cmd/synapbus/admin.go @@ -1441,16 +1441,18 @@ Examples: // ----- memory dream-run (feature 020 — manual dispatch) ----- var ( - dreamRunOwner string - dreamRunJobType string + dreamRunOwner string + dreamRunJobType string + dreamRunParallel int ) memoryDreamRunCmd := &cobra.Command{ Use: "dream-run", - Short: "Force a single consolidation job dispatch (bypasses trigger checks)", + Short: "Force a consolidation job dispatch (bypasses trigger checks). With --parallel N spawns N concurrent jobs.", RunE: func(cmd *cobra.Command, args []string) error { - resp, err := adminRequest("memory.dream_run", map[string]string{ + resp, err := adminRequest("memory.dream_run", map[string]any{ "owner": dreamRunOwner, "job_type": dreamRunJobType, + "parallel": dreamRunParallel, }) if err != nil { return err @@ -1461,6 +1463,7 @@ Examples: } memoryDreamRunCmd.Flags().StringVar(&dreamRunOwner, "owner", "", "Owner username or numeric user ID") memoryDreamRunCmd.Flags().StringVar(&dreamRunJobType, "job", "reflection", "Job type (reflection | core_rewrite | dedup_contradiction | link_gen)") + memoryDreamRunCmd.Flags().IntVar(&dreamRunParallel, "parallel", 0, "Spawn N concurrent jobs of this type (0 = use server SYNAPBUS_DREAM_PARALLEL; core_rewrite forces 1)") _ = memoryDreamRunCmd.MarkFlagRequired("owner") memoryCmd.AddCommand(memoryDreamRunCmd) diff --git a/cmd/synapbus/main.go b/cmd/synapbus/main.go index 4efc32a..1d2dbd1 100644 --- a/cmd/synapbus/main.go +++ b/cmd/synapbus/main.go @@ -512,7 +512,15 @@ func runServe(cmd *cobra.Command, args []string) error { // going through the existing createJob + poller path; subprocess // and webhook agents go through Registry.Execute. harnessRegistry := harness.NewRegistry() - harnessRegistry.Register(k8sjob.New(k8sRunner, nil, slog.Default())) + // Build a ClientsetWaiter when we have a real in-cluster runner so + // the k8sjob backend can actually wait for Job completion. Without + // this the backend errors immediately with "no Waiter configured" + // (the failure mode dream-worker jobs were hitting pre-fix). + var k8sWaiter k8sjob.Waiter + if rr, ok := k8sRunner.(*k8spkg.K8sJobRunner); ok { + k8sWaiter = k8sjob.NewClientsetWaiter(rr.GetClientset(), 0) + } + harnessRegistry.Register(k8sjob.New(k8sRunner, k8sWaiter, slog.Default())) // SYNAPBUS_KEEP_WORKDIR=1 preserves per-run workdirs after successful // runs. Useful when debugging MCP tool traces, gemini stdout, or // materialized config files. Default off to avoid disk growth. @@ -901,6 +909,10 @@ func runServe(cmd *cobra.Command, args []string) error { adminSvcs.DreamRun = func(ctx context.Context, ownerID, jobType string) (int64, error) { return c.ForceRun(ctx, ownerID, jobType) } + adminSvcs.DreamRunN = func(ctx context.Context, ownerID, jobType string, parallel int) ([]int64, error) { + return c.ForceRunN(ctx, ownerID, jobType, parallel) + } + adminSvcs.DefaultDreamParallel = memCfg.DreamParallel } adminServer := admin.NewServer(adminSocketPath, db.DB, adminSvcs, logger) if err := adminServer.Start(); err != nil { diff --git a/internal/admin/server.go b/internal/admin/server.go index a31fb12..27f962c 100644 --- a/internal/admin/server.go +++ b/internal/admin/server.go @@ -59,6 +59,15 @@ type Services struct { // 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. diff --git a/internal/admin/socket.go b/internal/admin/socket.go index c9ca6e8..d2b4e21 100644 --- a/internal/admin/socket.go +++ b/internal/admin/socket.go @@ -1987,13 +1987,15 @@ func (s *AdminServer) handleMemoryCoreDelete(ctx context.Context, args json.RawM }} } -// handleMemoryDreamRun forces a single dream-job dispatch. Bypasses -// trigger checks — useful for kubic verification (quickstart §"verify -// dream agent"). Returns the created job_id. +// handleMemoryDreamRun forces dream-job dispatch(es). With parallel=1 +// (default) returns one job_id. With parallel>1, fans out N concurrent +// jobs across slots 0..N-1 and returns the list of created ids. The +// circuit breaker still applies. func (s *AdminServer) handleMemoryDreamRun(ctx context.Context, args json.RawMessage) Response { var p struct { - Owner string `json:"owner"` - JobType string `json:"job_type"` + Owner string `json:"owner"` + JobType string `json:"job_type"` + Parallel int `json:"parallel"` } if err := json.Unmarshal(args, &p); err != nil { return Response{OK: false, Error: "invalid args: " + err.Error()} @@ -2001,22 +2003,34 @@ func (s *AdminServer) handleMemoryDreamRun(ctx context.Context, args json.RawMes if p.JobType == "" { return Response{OK: false, Error: "job_type is required"} } - if s.services.DreamRun == nil { + if s.services.DreamRunN == nil { return Response{OK: false, Error: "dream worker not configured (SYNAPBUS_DREAM_ENABLED=0?)"} } ownerStr, err := s.resolveOwnerString(ctx, p.Owner) if err != nil { return Response{OK: false, Error: err.Error()} } - jobID, err := s.services.DreamRun(ctx, ownerStr, p.JobType) - if err != nil { + parallel := p.Parallel + if parallel <= 0 { + parallel = s.services.DefaultDreamParallel + } + if parallel <= 0 { + parallel = 1 + } + ids, err := s.services.DreamRunN(ctx, ownerStr, p.JobType, parallel) + if err != nil && len(ids) == 0 { return Response{OK: false, Error: err.Error()} } - return Response{OK: true, Data: map[string]any{ - "job_id": jobID, + out := map[string]any{ + "job_ids": ids, "owner_id": ownerStr, "job_type": p.JobType, - }} + "parallel": parallel, + } + if err != nil { + out["error"] = err.Error() + } + return Response{OK: true, Data: out} } // Ensure the messaging import is used. diff --git a/internal/k8s/runner.go b/internal/k8s/runner.go index 5e5e019..8bb0989 100644 --- a/internal/k8s/runner.go +++ b/internal/k8s/runner.go @@ -2,11 +2,14 @@ package k8s import ( "context" + "crypto/rand" + "encoding/hex" "fmt" "io" "log/slog" "os" "strings" + "time" batchv1 "k8s.io/api/batch/v1" corev1 "k8s.io/api/core/v1" @@ -94,7 +97,17 @@ func (r *K8sJobRunner) GetNamespace() string { } func (r *K8sJobRunner) CreateJob(ctx context.Context, handler *K8sHandler, msg *JobMessage) (string, error) { - jobName := sanitizeJobName(fmt.Sprintf("synapbus-%s-%d", handler.AgentName, msg.MessageID)) + // When there is no triggering message (e.g. dream-worker dispatches + // where Message=nil → MessageID=0), fall back to a unique suffix so + // concurrent runs don't collide on the Job name. Format keeps the + // historic "synapbus--" prefix for log/grep continuity. + suffix := fmt.Sprintf("%d", msg.MessageID) + if msg.MessageID == 0 { + var b [4]byte + _, _ = rand.Read(b[:]) + suffix = fmt.Sprintf("%d-%s", time.Now().UnixNano()%1_000_000, hex.EncodeToString(b[:])) + } + jobName := sanitizeJobName(fmt.Sprintf("synapbus-%s-%s", handler.AgentName, suffix)) namespace := handler.Namespace if namespace == "" { diff --git a/internal/mcp/memory_tools.go b/internal/mcp/memory_tools.go index 067032d..515d1b7 100644 --- a/internal/mcp/memory_tools.go +++ b/internal/mcp/memory_tools.go @@ -284,6 +284,14 @@ func (r *MemoryToolRegistrar) handleListUnprocessed(ctx context.Context, req mcp windowExpr := fmt.Sprintf("-%d days", windowDays) queryArgs = append(queryArgs, owner, since, windowExpr, limit) + // Contract guarantees this list excludes: + // - messages already linked as the dst_message_id of a + // refines/duplicate_of/superseded_by edge (already + // consolidated by an earlier dream pass), AND + // - the dream worker's own output (from_agent prefix "dream:") + // so the agent never re-refines its own reflections. + // Without these filters the agent loops on the same oldest-50 + // messages every cycle and progress flat-lines. q := `SELECT m.id, m.from_agent, c.name, m.body, m.created_at FROM messages m JOIN agents a ON m.from_agent = a.name @@ -292,6 +300,11 @@ func (r *MemoryToolRegistrar) handleListUnprocessed(ctx context.Context, req mcp AND CAST(a.owner_id AS TEXT) = ? AND m.id > ? AND m.created_at > datetime('now', ?) + AND m.from_agent NOT LIKE 'dream:%' + AND m.id NOT IN ( + SELECT dst_message_id FROM memory_links + WHERE relation_type IN ('refines','duplicate_of','superseded_by') + ) ORDER BY m.id ASC LIMIT ?` diff --git a/internal/messaging/consolidation_jobs.go b/internal/messaging/consolidation_jobs.go index a49d5d5..d82b775 100644 --- a/internal/messaging/consolidation_jobs.go +++ b/internal/messaging/consolidation_jobs.go @@ -75,10 +75,18 @@ func NewJobsStore(db *sql.DB) *JobsStore { return &JobsStore{db: db} } -// Create inserts a `pending` row for the (owner, jobType) pair. If -// another job of the same type is already in flight, returns -// ErrJobAlreadyInFlight (mapped from the partial-unique-index conflict). +// Create inserts a `pending` row for the (owner, jobType) pair on slot 0. +// Kept for compatibility with single-threaded callers. Returns +// ErrJobAlreadyInFlight if slot 0 is occupied. func (s *JobsStore) Create(ctx context.Context, ownerID, jobType, triggerReason string) (int64, error) { + return s.CreateOnSlot(ctx, ownerID, jobType, triggerReason, 0) +} + +// CreateOnSlot inserts a `pending` row on a specific slot. Migration 030 +// extended the partial-unique index to include `slot`, so slots 0..N-1 +// can hold concurrent in-flight jobs of the same type per owner. The +// worker / admin CLI uses this when fanning out a backlog drain. +func (s *JobsStore) CreateOnSlot(ctx context.Context, ownerID, jobType, triggerReason string, slot int) (int64, error) { if s == nil || s.db == nil { return 0, fmt.Errorf("jobs store: nil store") } @@ -88,16 +96,16 @@ func (s *JobsStore) Create(ctx context.Context, ownerID, jobType, triggerReason if jobType == "" { return 0, fmt.Errorf("jobs store: empty job_type") } + if slot < 0 { + return 0, fmt.Errorf("jobs store: negative slot") + } res, err := s.db.ExecContext(ctx, `INSERT INTO memory_consolidation_jobs - (owner_id, job_type, status, trigger_reason) - VALUES (?, ?, 'pending', ?)`, - ownerID, jobType, triggerReason, + (owner_id, job_type, status, trigger_reason, slot) + VALUES (?, ?, 'pending', ?, ?)`, + ownerID, jobType, triggerReason, slot, ) if err != nil { - // modernc.org/sqlite surfaces unique-constraint conflicts via - // error strings; the partial-unique index is the only UNIQUE - // constraint that can fire here for INSERT. if isUniqueConstraint(err) { return 0, ErrJobAlreadyInFlight } @@ -107,6 +115,26 @@ func (s *JobsStore) Create(ctx context.Context, ownerID, jobType, triggerReason return id, nil } +// CreateNextAvailableSlot tries slots 0..maxSlots-1 in order and returns +// the (jobID, slot) of the first one that wasn't already in flight. Used +// when the worker / CLI wants to fan out N parallel jobs per cycle. +// Returns ErrJobAlreadyInFlight if all slots are busy. +func (s *JobsStore) CreateNextAvailableSlot(ctx context.Context, ownerID, jobType, triggerReason string, maxSlots int) (int64, int, error) { + if maxSlots <= 0 { + maxSlots = 1 + } + for slot := 0; slot < maxSlots; slot++ { + id, err := s.CreateOnSlot(ctx, ownerID, jobType, triggerReason, slot) + if err == nil { + return id, slot, nil + } + if !errors.Is(err, ErrJobAlreadyInFlight) { + return 0, 0, err + } + } + return 0, 0, ErrJobAlreadyInFlight +} + // Dispatch flips a pending row to `dispatched` and stamps the harness // run id and dispatch token. Returns an error if the row is not in // `pending` state. diff --git a/internal/messaging/consolidator.go b/internal/messaging/consolidator.go index 2ff6394..e65b6a7 100644 --- a/internal/messaging/consolidator.go +++ b/internal/messaging/consolidator.go @@ -359,12 +359,29 @@ func (w *ConsolidatorWorker) unprocessedCount(ctx context.Context, ownerID strin // ForceRun bypasses watermark/cron triggers and dispatches one job // for (ownerID, jobType) immediately. Used by the admin CLI // `synapbus memory dream-run` command. Returns the created job_id (or -// the existing in-flight one). +// the existing in-flight one). Uses slot 0 — for fan-out, see ForceRunN. // // The circuit breaker still applies — admins who want to override // must clear today's usage row directly. Manual override on top of a // blown budget defeats the safety net. func (w *ConsolidatorWorker) ForceRun(ctx context.Context, ownerID, jobType string) (int64, error) { + ids, err := w.ForceRunN(ctx, ownerID, jobType, 1) + if err != nil || len(ids) == 0 { + return 0, err + } + return ids[0], nil +} + +// ForceRunN dispatches up to N parallel jobs of the same type for one +// owner. core_rewrite ignores N and dispatches at most one (per-agent +// blob is wholesale-replace; concurrent rewrites would race). +func (w *ConsolidatorWorker) ForceRunN(ctx context.Context, ownerID, jobType string, parallel int) ([]int64, error) { + if parallel <= 0 { + parallel = 1 + } + if jobType == JobTypeCoreRewrite { + parallel = 1 + } if w.gate != nil { allowed, reason, _ := w.gate.Allow(ctx, ownerID) if !allowed { @@ -376,37 +393,58 @@ func (w *ConsolidatorWorker) ForceRun(ctx context.Context, ownerID, jobType stri } recordCircuitBrokenMetric(ownerID, jobType, reason) recordJobMetric(ownerID, jobType, JobStatusCircuitBroken) - return jobID, fmt.Errorf("circuit broken: %s", reason) + return []int64{jobID}, fmt.Errorf("circuit broken: %s", reason) } - return 0, fmt.Errorf("circuit broken: %s", reason) + return nil, fmt.Errorf("circuit broken: %s", reason) } } - jobID, err := w.jobs.Create(ctx, ownerID, jobType, "manual:"+ownerID) - if err != nil { - if errors.Is(err, ErrJobAlreadyInFlight) { - if active, _ := w.jobs.ActiveJob(ctx, ownerID, jobType); active != nil { - return active.ID, nil + out := make([]int64, 0, parallel) + for slot := 0; slot < parallel; slot++ { + jobID, err := w.jobs.CreateOnSlot(ctx, ownerID, jobType, "manual:"+ownerID, slot) + if err != nil { + if errors.Is(err, ErrJobAlreadyInFlight) { + // Slot already busy. If this is the first slot, fall back + // to returning the existing in-flight job (preserves + // historical ForceRun semantics). + if len(out) == 0 && slot == 0 { + if active, _ := w.jobs.ActiveJob(ctx, ownerID, jobType); active != nil { + return []int64{active.ID}, nil + } + } + w.logger.Debug("slot busy; skipping", "owner_id", ownerID, "job_type", jobType, "slot", slot) + continue } + return out, err } - return 0, err + if err := w.launchOne(ctx, ownerID, jobType, jobID); err != nil { + return out, err + } + out = append(out, jobID) } + return out, nil +} + +// launchOne wires up the per-job tokens / harness dispatch for an +// already-Created job row. Shared by ForceRunN and the worker's +// internal tryDispatch path. +func (w *ConsolidatorWorker) launchOne(ctx context.Context, ownerID, jobType string, jobID int64) error { if w.usage != nil { _ = w.usage.RecordStart(ctx, ownerID) } tok, _, err := w.tokens.Issue(ctx, ownerID, jobID) if err != nil { _ = w.jobs.Complete(ctx, jobID, JobStatusFailed, "", "token issue: "+err.Error()) - return 0, fmt.Errorf("issue token: %w", err) + return fmt.Errorf("issue token: %w", err) } agent, err := w.agentLook.GetAgent(ctx, w.cfg.DreamAgent) if err != nil || agent == nil { _ = w.jobs.Complete(ctx, jobID, JobStatusFailed, "", "dream agent not found") - return 0, fmt.Errorf("dream agent %q not found: %w", w.cfg.DreamAgent, err) + return fmt.Errorf("dream agent %q not found: %w", w.cfg.DreamAgent, err) } runID := uuid.NewString() if err := w.jobs.Dispatch(ctx, jobID, runID, tok); err != nil { _ = w.jobs.Complete(ctx, jobID, JobStatusFailed, "", "dispatch flip: "+err.Error()) - return 0, fmt.Errorf("dispatch flip: %w", err) + return fmt.Errorf("dispatch flip: %w", err) } w.wg.Add(1) go func() { @@ -419,7 +457,7 @@ func (w *ConsolidatorWorker) ForceRun(ctx context.Context, ownerID, jobType stri } w.runJob(ownerID, jobID, jobType, tok, runID, agent) }() - return jobID, nil + return nil } // tryDispatch attempts to create+dispatch one job. Idempotent — diff --git a/internal/messaging/memory_config.go b/internal/messaging/memory_config.go index 80e777f..b030b46 100644 --- a/internal/messaging/memory_config.go +++ b/internal/messaging/memory_config.go @@ -45,6 +45,14 @@ type MemoryConfig struct { // DreamMaxConcurrent caps the number of owners with in-flight jobs. DreamMaxConcurrent int + // DreamParallel is the number of concurrent dream jobs the worker + // (and `synapbus memory dream-run`) will fan out per (owner, job_type) + // when triggered. 1 = historical single-job behaviour. Higher values + // let one watermark/manual trigger spawn N parallel reflection / + // link_gen / dedup_contradiction agents to drain a backlog. core_rewrite + // is always single-slot regardless of this knob. + DreamParallel int + // DreamWallclockBudget caps how long a single consolidation job is // allowed to run before the worker terminates it with status=partial. DreamWallclockBudget time.Duration @@ -96,6 +104,7 @@ func DefaultMemoryConfig() MemoryConfig { DreamInterval: 1 * time.Hour, DreamDeepCron: "0 3 * * *", DreamMaxConcurrent: 4, + DreamParallel: 1, DreamWallclockBudget: 10 * time.Minute, DreamWatermark: 20, DreamAgent: "claude-code", @@ -156,6 +165,11 @@ func ParseMemoryConfig() MemoryConfig { cfg.DreamMaxConcurrent = n } } + if v := os.Getenv("SYNAPBUS_DREAM_PARALLEL"); v != "" { + if n, err := strconv.Atoi(v); err == nil && n > 0 { + cfg.DreamParallel = n + } + } if v := os.Getenv("SYNAPBUS_DREAM_WALLCLOCK_BUDGET"); v != "" { if d, err := parseDurationWithDays(v); err == nil && d > 0 { cfg.DreamWallclockBudget = d diff --git a/internal/storage/schema/030_dream_parallelism.sql b/internal/storage/schema/030_dream_parallelism.sql new file mode 100644 index 0000000..72446b2 --- /dev/null +++ b/internal/storage/schema/030_dream_parallelism.sql @@ -0,0 +1,15 @@ +-- 030_dream_parallelism.sql +-- +-- Allow concurrent dream-worker dispatches of the same (owner, job_type) +-- by adding a `slot` discriminator. Slot 0 is the historical behaviour +-- (one in-flight per type). Slots 1..N-1 are used when the worker / CLI +-- fans out to drain a backlog quickly. + +ALTER TABLE memory_consolidation_jobs + ADD COLUMN slot INTEGER NOT NULL DEFAULT 0; + +DROP INDEX IF EXISTS idx_consolidation_in_flight; + +CREATE UNIQUE INDEX IF NOT EXISTS idx_consolidation_in_flight + ON memory_consolidation_jobs(owner_id, job_type, slot) + WHERE status IN ('pending', 'dispatched', 'running');