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) -----
|
// ----- memory dream-run (feature 020 — manual dispatch) -----
|
||||||
var (
|
var (
|
||||||
dreamRunOwner string
|
dreamRunOwner string
|
||||||
dreamRunJobType string
|
dreamRunJobType string
|
||||||
|
dreamRunParallel int
|
||||||
)
|
)
|
||||||
memoryDreamRunCmd := &cobra.Command{
|
memoryDreamRunCmd := &cobra.Command{
|
||||||
Use: "dream-run",
|
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 {
|
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,
|
"owner": dreamRunOwner,
|
||||||
"job_type": dreamRunJobType,
|
"job_type": dreamRunJobType,
|
||||||
|
"parallel": dreamRunParallel,
|
||||||
})
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
@@ -1461,6 +1463,7 @@ Examples:
|
|||||||
}
|
}
|
||||||
memoryDreamRunCmd.Flags().StringVar(&dreamRunOwner, "owner", "", "Owner username or numeric user ID")
|
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().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")
|
_ = memoryDreamRunCmd.MarkFlagRequired("owner")
|
||||||
memoryCmd.AddCommand(memoryDreamRunCmd)
|
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
|
// going through the existing createJob + poller path; subprocess
|
||||||
// and webhook agents go through Registry.Execute.
|
// and webhook agents go through Registry.Execute.
|
||||||
harnessRegistry := harness.NewRegistry()
|
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
|
// SYNAPBUS_KEEP_WORKDIR=1 preserves per-run workdirs after successful
|
||||||
// runs. Useful when debugging MCP tool traces, gemini stdout, or
|
// runs. Useful when debugging MCP tool traces, gemini stdout, or
|
||||||
// materialized config files. Default off to avoid disk growth.
|
// 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) {
|
adminSvcs.DreamRun = func(ctx context.Context, ownerID, jobType string) (int64, error) {
|
||||||
return c.ForceRun(ctx, ownerID, jobType)
|
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)
|
adminServer := admin.NewServer(adminSocketPath, db.DB, adminSvcs, logger)
|
||||||
if err := adminServer.Start(); err != nil {
|
if err := adminServer.Start(); err != nil {
|
||||||
|
|||||||
@@ -59,6 +59,15 @@ type Services struct {
|
|||||||
// consolidator worker is enabled. Closure form keeps the worker
|
// consolidator worker is enabled. Closure form keeps the worker
|
||||||
// internals out of the admin package's import graph.
|
// internals out of the admin package's import graph.
|
||||||
DreamRun func(ctx context.Context, ownerID, jobType string) (jobID int64, err error)
|
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.
|
// 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
|
// handleMemoryDreamRun forces dream-job dispatch(es). With parallel=1
|
||||||
// trigger checks — useful for kubic verification (quickstart §"verify
|
// (default) returns one job_id. With parallel>1, fans out N concurrent
|
||||||
// dream agent"). Returns the created job_id.
|
// 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 {
|
func (s *AdminServer) handleMemoryDreamRun(ctx context.Context, args json.RawMessage) Response {
|
||||||
var p struct {
|
var p struct {
|
||||||
Owner string `json:"owner"`
|
Owner string `json:"owner"`
|
||||||
JobType string `json:"job_type"`
|
JobType string `json:"job_type"`
|
||||||
|
Parallel int `json:"parallel"`
|
||||||
}
|
}
|
||||||
if err := json.Unmarshal(args, &p); err != nil {
|
if err := json.Unmarshal(args, &p); err != nil {
|
||||||
return Response{OK: false, Error: "invalid args: " + err.Error()}
|
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 == "" {
|
if p.JobType == "" {
|
||||||
return Response{OK: false, Error: "job_type is required"}
|
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?)"}
|
return Response{OK: false, Error: "dream worker not configured (SYNAPBUS_DREAM_ENABLED=0?)"}
|
||||||
}
|
}
|
||||||
ownerStr, err := s.resolveOwnerString(ctx, p.Owner)
|
ownerStr, err := s.resolveOwnerString(ctx, p.Owner)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return Response{OK: false, Error: err.Error()}
|
return Response{OK: false, Error: err.Error()}
|
||||||
}
|
}
|
||||||
jobID, err := s.services.DreamRun(ctx, ownerStr, p.JobType)
|
parallel := p.Parallel
|
||||||
if err != nil {
|
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: false, Error: err.Error()}
|
||||||
}
|
}
|
||||||
return Response{OK: true, Data: map[string]any{
|
out := map[string]any{
|
||||||
"job_id": jobID,
|
"job_ids": ids,
|
||||||
"owner_id": ownerStr,
|
"owner_id": ownerStr,
|
||||||
"job_type": p.JobType,
|
"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.
|
// Ensure the messaging import is used.
|
||||||
|
|||||||
+14
-1
@@ -2,11 +2,14 @@ package k8s
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"crypto/rand"
|
||||||
|
"encoding/hex"
|
||||||
"fmt"
|
"fmt"
|
||||||
"io"
|
"io"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"os"
|
"os"
|
||||||
"strings"
|
"strings"
|
||||||
|
"time"
|
||||||
|
|
||||||
batchv1 "k8s.io/api/batch/v1"
|
batchv1 "k8s.io/api/batch/v1"
|
||||||
corev1 "k8s.io/api/core/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) {
|
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
|
namespace := handler.Namespace
|
||||||
if namespace == "" {
|
if namespace == "" {
|
||||||
|
|||||||
@@ -284,6 +284,14 @@ func (r *MemoryToolRegistrar) handleListUnprocessed(ctx context.Context, req mcp
|
|||||||
windowExpr := fmt.Sprintf("-%d days", windowDays)
|
windowExpr := fmt.Sprintf("-%d days", windowDays)
|
||||||
queryArgs = append(queryArgs, owner, since, windowExpr, limit)
|
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
|
q := `SELECT m.id, m.from_agent, c.name, m.body, m.created_at
|
||||||
FROM messages m
|
FROM messages m
|
||||||
JOIN agents a ON m.from_agent = a.name
|
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 CAST(a.owner_id AS TEXT) = ?
|
||||||
AND m.id > ?
|
AND m.id > ?
|
||||||
AND m.created_at > datetime('now', ?)
|
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
|
ORDER BY m.id ASC
|
||||||
LIMIT ?`
|
LIMIT ?`
|
||||||
|
|
||||||
|
|||||||
@@ -75,10 +75,18 @@ func NewJobsStore(db *sql.DB) *JobsStore {
|
|||||||
return &JobsStore{db: db}
|
return &JobsStore{db: db}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Create inserts a `pending` row for the (owner, jobType) pair. If
|
// Create inserts a `pending` row for the (owner, jobType) pair on slot 0.
|
||||||
// another job of the same type is already in flight, returns
|
// Kept for compatibility with single-threaded callers. Returns
|
||||||
// ErrJobAlreadyInFlight (mapped from the partial-unique-index conflict).
|
// ErrJobAlreadyInFlight if slot 0 is occupied.
|
||||||
func (s *JobsStore) Create(ctx context.Context, ownerID, jobType, triggerReason string) (int64, error) {
|
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 {
|
if s == nil || s.db == nil {
|
||||||
return 0, fmt.Errorf("jobs store: nil store")
|
return 0, fmt.Errorf("jobs store: nil store")
|
||||||
}
|
}
|
||||||
@@ -88,16 +96,16 @@ func (s *JobsStore) Create(ctx context.Context, ownerID, jobType, triggerReason
|
|||||||
if jobType == "" {
|
if jobType == "" {
|
||||||
return 0, fmt.Errorf("jobs store: empty job_type")
|
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,
|
res, err := s.db.ExecContext(ctx,
|
||||||
`INSERT INTO memory_consolidation_jobs
|
`INSERT INTO memory_consolidation_jobs
|
||||||
(owner_id, job_type, status, trigger_reason)
|
(owner_id, job_type, status, trigger_reason, slot)
|
||||||
VALUES (?, ?, 'pending', ?)`,
|
VALUES (?, ?, 'pending', ?, ?)`,
|
||||||
ownerID, jobType, triggerReason,
|
ownerID, jobType, triggerReason, slot,
|
||||||
)
|
)
|
||||||
if err != nil {
|
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) {
|
if isUniqueConstraint(err) {
|
||||||
return 0, ErrJobAlreadyInFlight
|
return 0, ErrJobAlreadyInFlight
|
||||||
}
|
}
|
||||||
@@ -107,6 +115,26 @@ func (s *JobsStore) Create(ctx context.Context, ownerID, jobType, triggerReason
|
|||||||
return id, nil
|
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
|
// 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
|
// run id and dispatch token. Returns an error if the row is not in
|
||||||
// `pending` state.
|
// `pending` state.
|
||||||
|
|||||||
@@ -359,12 +359,29 @@ func (w *ConsolidatorWorker) unprocessedCount(ctx context.Context, ownerID strin
|
|||||||
// ForceRun bypasses watermark/cron triggers and dispatches one job
|
// ForceRun bypasses watermark/cron triggers and dispatches one job
|
||||||
// for (ownerID, jobType) immediately. Used by the admin CLI
|
// for (ownerID, jobType) immediately. Used by the admin CLI
|
||||||
// `synapbus memory dream-run` command. Returns the created job_id (or
|
// `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
|
// The circuit breaker still applies — admins who want to override
|
||||||
// must clear today's usage row directly. Manual override on top of a
|
// must clear today's usage row directly. Manual override on top of a
|
||||||
// blown budget defeats the safety net.
|
// blown budget defeats the safety net.
|
||||||
func (w *ConsolidatorWorker) ForceRun(ctx context.Context, ownerID, jobType string) (int64, error) {
|
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 {
|
if w.gate != nil {
|
||||||
allowed, reason, _ := w.gate.Allow(ctx, ownerID)
|
allowed, reason, _ := w.gate.Allow(ctx, ownerID)
|
||||||
if !allowed {
|
if !allowed {
|
||||||
@@ -376,37 +393,58 @@ func (w *ConsolidatorWorker) ForceRun(ctx context.Context, ownerID, jobType stri
|
|||||||
}
|
}
|
||||||
recordCircuitBrokenMetric(ownerID, jobType, reason)
|
recordCircuitBrokenMetric(ownerID, jobType, reason)
|
||||||
recordJobMetric(ownerID, jobType, JobStatusCircuitBroken)
|
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)
|
out := make([]int64, 0, parallel)
|
||||||
if err != nil {
|
for slot := 0; slot < parallel; slot++ {
|
||||||
if errors.Is(err, ErrJobAlreadyInFlight) {
|
jobID, err := w.jobs.CreateOnSlot(ctx, ownerID, jobType, "manual:"+ownerID, slot)
|
||||||
if active, _ := w.jobs.ActiveJob(ctx, ownerID, jobType); active != nil {
|
if err != nil {
|
||||||
return active.ID, 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 {
|
if w.usage != nil {
|
||||||
_ = w.usage.RecordStart(ctx, ownerID)
|
_ = w.usage.RecordStart(ctx, ownerID)
|
||||||
}
|
}
|
||||||
tok, _, err := w.tokens.Issue(ctx, ownerID, jobID)
|
tok, _, err := w.tokens.Issue(ctx, ownerID, jobID)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
_ = w.jobs.Complete(ctx, jobID, JobStatusFailed, "", "token issue: "+err.Error())
|
_ = 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)
|
agent, err := w.agentLook.GetAgent(ctx, w.cfg.DreamAgent)
|
||||||
if err != nil || agent == nil {
|
if err != nil || agent == nil {
|
||||||
_ = w.jobs.Complete(ctx, jobID, JobStatusFailed, "", "dream agent not found")
|
_ = 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()
|
runID := uuid.NewString()
|
||||||
if err := w.jobs.Dispatch(ctx, jobID, runID, tok); err != nil {
|
if err := w.jobs.Dispatch(ctx, jobID, runID, tok); err != nil {
|
||||||
_ = w.jobs.Complete(ctx, jobID, JobStatusFailed, "", "dispatch flip: "+err.Error())
|
_ = 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)
|
w.wg.Add(1)
|
||||||
go func() {
|
go func() {
|
||||||
@@ -419,7 +457,7 @@ func (w *ConsolidatorWorker) ForceRun(ctx context.Context, ownerID, jobType stri
|
|||||||
}
|
}
|
||||||
w.runJob(ownerID, jobID, jobType, tok, runID, agent)
|
w.runJob(ownerID, jobID, jobType, tok, runID, agent)
|
||||||
}()
|
}()
|
||||||
return jobID, nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// tryDispatch attempts to create+dispatch one job. Idempotent —
|
// 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 caps the number of owners with in-flight jobs.
|
||||||
DreamMaxConcurrent int
|
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
|
// DreamWallclockBudget caps how long a single consolidation job is
|
||||||
// allowed to run before the worker terminates it with status=partial.
|
// allowed to run before the worker terminates it with status=partial.
|
||||||
DreamWallclockBudget time.Duration
|
DreamWallclockBudget time.Duration
|
||||||
@@ -96,6 +104,7 @@ func DefaultMemoryConfig() MemoryConfig {
|
|||||||
DreamInterval: 1 * time.Hour,
|
DreamInterval: 1 * time.Hour,
|
||||||
DreamDeepCron: "0 3 * * *",
|
DreamDeepCron: "0 3 * * *",
|
||||||
DreamMaxConcurrent: 4,
|
DreamMaxConcurrent: 4,
|
||||||
|
DreamParallel: 1,
|
||||||
DreamWallclockBudget: 10 * time.Minute,
|
DreamWallclockBudget: 10 * time.Minute,
|
||||||
DreamWatermark: 20,
|
DreamWatermark: 20,
|
||||||
DreamAgent: "claude-code",
|
DreamAgent: "claude-code",
|
||||||
@@ -156,6 +165,11 @@ func ParseMemoryConfig() MemoryConfig {
|
|||||||
cfg.DreamMaxConcurrent = n
|
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 v := os.Getenv("SYNAPBUS_DREAM_WALLCLOCK_BUDGET"); v != "" {
|
||||||
if d, err := parseDurationWithDays(v); err == nil && d > 0 {
|
if d, err := parseDurationWithDays(v); err == nil && d > 0 {
|
||||||
cfg.DreamWallclockBudget = d
|
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