Files
Algis DumbrisandClaude Opus 4.7 bfb2551b45 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>
2026-05-12 19:53:25 +03:00

302 lines
8.4 KiB
Go

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"
"k8s.io/apimachinery/pkg/api/resource"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/rest"
)
// JobRunner is the interface for creating and managing K8s Jobs.
type JobRunner interface {
// IsAvailable returns true if the K8s runner is available (running in-cluster).
IsAvailable() bool
// CreateJob creates a K8s Job for the given handler and message.
CreateJob(ctx context.Context, handler *K8sHandler, msg *JobMessage) (string, error)
// GetJobLogs returns the logs for a completed Job.
GetJobLogs(ctx context.Context, namespace, jobName string) (string, error)
// GetNamespace returns the namespace SynapBus is running in.
GetNamespace() string
}
// JobMessage contains the message data to inject into the K8s Job.
type JobMessage struct {
MessageID int64
FromAgent string
Body string
Event string
Channel string
Timestamp string
}
// K8sJobRunner implements JobRunner using the K8s API.
type K8sJobRunner struct {
clientset kubernetes.Interface
namespace string
logger *slog.Logger
}
// NewJobRunner attempts to create a K8s runner using in-cluster config.
// If not running in a K8s cluster, returns a NoopRunner.
func NewJobRunner(logger *slog.Logger) JobRunner {
config, err := rest.InClusterConfig()
if err != nil {
logger.Info("not running in Kubernetes cluster, K8s job runner disabled", "reason", err.Error())
return NewNoopRunner()
}
clientset, err := kubernetes.NewForConfig(config)
if err != nil {
logger.Error("failed to create K8s client", "error", err)
return NewNoopRunner()
}
// Detect current namespace
ns := detectNamespace()
logger.Info("Kubernetes job runner initialized", "namespace", ns)
return &K8sJobRunner{
clientset: clientset,
namespace: ns,
logger: logger,
}
}
// NewJobRunnerWithClient creates a K8s runner with a provided clientset (for testing).
func NewJobRunnerWithClient(clientset kubernetes.Interface, namespace string, logger *slog.Logger) *K8sJobRunner {
return &K8sJobRunner{
clientset: clientset,
namespace: namespace,
logger: logger,
}
}
func (r *K8sJobRunner) IsAvailable() bool {
return true
}
// GetClientset returns the kubernetes clientset for direct API access (used by reactor poller).
func (r *K8sJobRunner) GetClientset() kubernetes.Interface {
return r.clientset
}
func (r *K8sJobRunner) GetNamespace() string {
return r.namespace
}
func (r *K8sJobRunner) CreateJob(ctx context.Context, handler *K8sHandler, msg *JobMessage) (string, error) {
// 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 == "" {
namespace = r.namespace
}
// Build environment variables
envVars := []corev1.EnvVar{
{Name: "SYNAPBUS_MESSAGE_ID", Value: fmt.Sprintf("%d", msg.MessageID)},
{Name: "SYNAPBUS_MESSAGE_BODY", Value: truncateBody(msg.Body, 32768)},
{Name: "SYNAPBUS_FROM_AGENT", Value: msg.FromAgent},
{Name: "SYNAPBUS_EVENT", Value: msg.Event},
{Name: "SYNAPBUS_TIMESTAMP", Value: msg.Timestamp},
}
if msg.Channel != "" {
envVars = append(envVars, corev1.EnvVar{Name: "SYNAPBUS_CHANNEL", Value: msg.Channel})
}
// Add user-defined env vars
for k, v := range handler.Env {
envVars = append(envVars, corev1.EnvVar{Name: k, Value: v})
}
// Build resource limits
resourceLimits := corev1.ResourceList{}
if handler.ResourcesMemory != "" {
resourceLimits[corev1.ResourceMemory] = resource.MustParse(handler.ResourcesMemory)
}
if handler.ResourcesCPU != "" {
resourceLimits[corev1.ResourceCPU] = resource.MustParse(handler.ResourcesCPU)
}
backoffLimit := int32(0)
activeDeadline := int64(handler.TimeoutSeconds)
ttlAfterFinished := int32(3600)
job := &batchv1.Job{
ObjectMeta: metav1.ObjectMeta{
Name: jobName,
Namespace: namespace,
Labels: map[string]string{
"app.kubernetes.io/managed-by": "synapbus",
"synapbus.io/agent": handler.AgentName,
"synapbus.io/handler-id": fmt.Sprintf("%d", handler.ID),
},
},
Spec: batchv1.JobSpec{
BackoffLimit: &backoffLimit,
ActiveDeadlineSeconds: &activeDeadline,
TTLSecondsAfterFinished: &ttlAfterFinished,
Template: corev1.PodTemplateSpec{
Spec: corev1.PodSpec{
RestartPolicy: corev1.RestartPolicyNever,
Containers: []corev1.Container{
{
Name: "handler",
Image: handler.Image,
ImagePullPolicy: corev1.PullIfNotPresent,
Args: handler.Args,
Env: envVars,
VolumeMounts: buildVolumeMounts(handler.VolumeMounts),
Resources: corev1.ResourceRequirements{
Limits: resourceLimits,
},
},
},
Volumes: buildVolumes(handler.Volumes),
},
},
},
}
created, err := r.clientset.BatchV1().Jobs(namespace).Create(ctx, job, metav1.CreateOptions{})
if err != nil {
return "", fmt.Errorf("create K8s Job: %w", err)
}
r.logger.Info("K8s Job created",
"job_name", created.Name,
"namespace", namespace,
"agent", handler.AgentName,
"image", handler.Image,
)
return created.Name, nil
}
func (r *K8sJobRunner) GetJobLogs(ctx context.Context, namespace, jobName string) (string, error) {
// Find pods for this job
pods, err := r.clientset.CoreV1().Pods(namespace).List(ctx, metav1.ListOptions{
LabelSelector: fmt.Sprintf("job-name=%s", jobName),
})
if err != nil {
return "", fmt.Errorf("list pods for job %s: %w", jobName, err)
}
if len(pods.Items) == 0 {
return "", fmt.Errorf("no pods found for job %s", jobName)
}
// Get logs from the first pod
pod := pods.Items[0]
logStream, err := r.clientset.CoreV1().Pods(namespace).GetLogs(pod.Name, &corev1.PodLogOptions{}).Stream(ctx)
if err != nil {
return "", fmt.Errorf("get logs for pod %s: %w", pod.Name, err)
}
defer logStream.Close()
logs, err := io.ReadAll(io.LimitReader(logStream, 1<<20)) // 1MB limit
if err != nil {
return "", fmt.Errorf("read logs: %w", err)
}
return string(logs), nil
}
// detectNamespace reads the current namespace from the mounted service account.
func detectNamespace() string {
data, err := os.ReadFile("/var/run/secrets/kubernetes.io/serviceaccount/namespace")
if err == nil && len(data) > 0 {
return strings.TrimSpace(string(data))
}
return "default"
}
// sanitizeJobName ensures the job name is valid for Kubernetes (lowercase, max 63 chars, DNS-safe).
func sanitizeJobName(name string) string {
name = strings.ToLower(name)
name = strings.Map(func(r rune) rune {
if (r >= 'a' && r <= 'z') || (r >= '0' && r <= '9') || r == '-' {
return r
}
return '-'
}, name)
// Trim leading/trailing hyphens
name = strings.Trim(name, "-")
if len(name) > 63 {
name = name[:63]
}
return name
}
// buildVolumeMounts converts our VolumeMount type to K8s VolumeMounts.
func buildVolumeMounts(mounts []VolumeMount) []corev1.VolumeMount {
if len(mounts) == 0 {
return nil
}
var result []corev1.VolumeMount
for _, m := range mounts {
result = append(result, corev1.VolumeMount{
Name: m.Name,
MountPath: m.MountPath,
ReadOnly: m.ReadOnly,
})
}
return result
}
// buildVolumes converts our Volume type to K8s Volumes.
func buildVolumes(volumes []Volume) []corev1.Volume {
if len(volumes) == 0 {
return nil
}
var result []corev1.Volume
for _, v := range volumes {
vol := corev1.Volume{Name: v.Name}
if v.HostPath != "" {
hostPathType := corev1.HostPathDirectory
vol.VolumeSource = corev1.VolumeSource{
HostPath: &corev1.HostPathVolumeSource{
Path: v.HostPath,
Type: &hostPathType,
},
}
} else if v.EmptyDir {
vol.VolumeSource = corev1.VolumeSource{
EmptyDir: &corev1.EmptyDirVolumeSource{},
}
}
result = append(result, vol)
}
return result
}
// truncateBody truncates the message body to maxLen bytes.
func truncateBody(body string, maxLen int) string {
if len(body) <= maxLen {
return body
}
return body[:maxLen]
}