feat(020): configurable dream parallelism + 3 bug fixes from kubic drain

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>
This commit is contained in:
Algis Dumbris
2026-05-12 19:53:25 +03:00
co-authored by Claude Opus 4.7
parent 1d894ec5aa
commit bfb2551b45
10 changed files with 198 additions and 39 deletions
+7 -4
View File
@@ -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)
+13 -1
View File
@@ -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 {
+9
View File
@@ -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.
+25 -11
View File
@@ -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.
+14 -1
View File
@@ -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-<agent>-<id>" 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 == "" {
+13
View File
@@ -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 ?`
+37 -9
View File
@@ -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.
+51 -13
View File
@@ -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 —
+14
View File
@@ -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
@@ -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');