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:
co-authored by
Claude Opus 4.7
parent
1d894ec5aa
commit
bfb2551b45
@@ -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
@@ -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 {
|
||||
|
||||
@@ -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
@@ -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
@@ -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 == "" {
|
||||
|
||||
@@ -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 ?`
|
||||
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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 —
|
||||
|
||||
@@ -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');
|
||||
Reference in New Issue
Block a user