dreamer

package
v0.3.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 7, 2026 License: Apache-2.0 Imports: 34 Imported by: 0

Documentation

Overview

Package dreamer is the mining loop's dreaming: the service that keeps the slice dispatcher's ready frontier as long as its admission can use, with the work most likely to bend the struggle rate down. It generates and ranks work; it never executes it, and it never launches a session.

The pass

Dreamer.Dream is one pass, which the binary runs as a cron trigger every DreamInterval, after every merge, and once on start. A pass reads the dispatcher's snapshot and slice graph, releases the slices it added on the pass before (each is added held, so the Workbench queue shows it with its source and rank before the dispatcher can launch it), and derives the target: the launches the admission allows now plus the sessions expected to finish before the next pass, from the measured enqueue-to-merge durations. When the ready frontier is shorter, it asks every source for candidates, ranks them and adds the best through the dispatcher's AddSlice until it is not.

It adds nothing while it is paused, while the dispatcher is paused, or while any of the dispatcher's limits (the daily budget, the harness's load and disk floors, the provider's rate) admits no launch. An operator's hold is never released: the dreamer releases only its own visibility hold.

Ranking

A source is a named function value registered with WithSource. A candidate carries how often the problem it addresses occurs per week and the measured outcomes of earlier work on the same class. Its score is frequency × fixability × exploration: fixability is the Laplace estimate (fixed + 1) / (attempted + 2) over those outcomes, pooled with the outcomes of the dreamer's own slices of that source; exploration is the strategy's weight on classes nothing has worked on before, and 1 for known ones. Fallback sources (the harness's own gaps) are asked only when every primary source found nothing.

Strategy and the held-out sessions

HeldOut splits the corpus's runs deterministically. The struggle-class source chooses work from the in-sample runs only; Compare measures the weekly struggle factor on both. Dreamer.Strategize, run daily, sets the exploration weight from the held-out series alone: one plus the number of trailing complete weeks whose held-out rate did not fall.

Every pass, control and strategy is a typed Record, handed to the binary's record sink; the latest state is a Snapshot the Workbench shows.

Index

Constants

View Source
const (
	// TriggerDream is the dreamer's pass, every DreamInterval.
	TriggerDream = "dreamer.dream"
	// TriggerStrategy recomputes the held-out comparison and the strategy
	// once a day; the strategy changes at most once a week, when a week
	// completes.
	TriggerStrategy = "dreamer.strategy"
	// DreamInterval is the cadence of the passes: a slice the dreamer adds is
	// on the Workbench queue for one interval before the next pass releases
	// it to the dispatcher.
	DreamInterval = 3 * time.Minute

	// SnapshotFile is where the binary writes every published snapshot, under
	// the harness state directory, for the Workbench.
	SnapshotFile = "dreamer.json"
	// RecordFile is where the binary appends every record.
	RecordFile = "dreamer.jsonl"
	// VisibilityHold is the reason the dreamer holds every slice it adds. The
	// next pass releases a slice held with exactly this reason; an operator's
	// hold replaces the reason, and the dreamer never releases it.
	VisibilityHold = "dreamer: shown on the Workbench queue for one dreamer pass before the dispatcher may launch it"
	// ProvenancePrefix opens the provenance source of every slice the
	// dreamer adds.
	ProvenancePrefix = "dreamer"
)
View Source
const (
	SnapshotTool   = "DreamerSnapshot"
	DreamTool      = "Dream"
	PauseTool      = "PauseDreamer"
	ResumeTool     = "ResumeDreamer"
	StrategizeTool = "DreamerStrategy"

	SnapshotPath   = "/api/dreamer"
	DreamPath      = "/api/dreamer/dream"
	PausePath      = "/api/dreamer/pause"
	ResumePath     = "/api/dreamer/resume"
	StrategizePath = "/api/dreamer/strategy"
)

The dreamer's typed operations: each is one MCP tool on the CSF service and one HTTP route.

View Source
const (
	MetricReady     = "csf_dreamer_ready_slices"
	MetricTarget    = "csf_dreamer_target_slices"
	MetricGenerated = "csf_dreamer_slices_generated"
	MetricFactor    = "csf_dreamer_struggle_factor_per_week"
	MetricRate      = "csf_dreamer_struggle_rate_per_1k_tool_calls"
	MetricExplore   = "csf_dreamer_explore_weight"
)

The families the dreamer exports, each a panel of dashboard.json.

View Source
const (
	PartitionAll      = "all"
	PartitionInSample = "in_sample"
	PartitionHeldOut  = "held_out"
)

The partitions a comparison reports.

Variables

View Source
var (
	ErrInvalidOption  = errors.New("dreamer: invalid option")
	ErrNoDispatcher   = errors.New("dreamer: a dispatcher is required")
	ErrNoRecipe       = errors.New("dreamer: a recipe with a workspace is required")
	ErrNotStarted     = errors.New("dreamer: not started")
	ErrInvalidControl = errors.New("dreamer: invalid control")
)

The errors the service reports.

Functions

func Brief

func Brief(entry Ranked, total int) string

Brief is the task a generated slice's session runs after the recipe's own: what it is, why it ranks where it does, its evidence and its acceptance.

func HeldOut

func HeldOut(run string) bool

HeldOut reports whether a run is held out: the FNV-1a hash of its name modulo heldOutShare is zero. It is deterministic, so a run never changes sides, and the dreamer never reads a held-out run to choose work.

func RunOf

func RunOf(item string) string

RunOf is the run an item of the corpus belongs to: its first path element.

func SliceID

func SliceID(candidate Candidate) string

SliceID is the slice a candidate becomes: its source and a digest of its key, so the same work is never added twice.

Types

type Added

type Added struct {
	SliceID   string     `json:"slice_id"`
	Source    SourceName `json:"source"`
	Key       string     `json:"key"`
	Rank      int        `json:"rank"`
	Score     float64    `json:"score"`
	TicketURL string     `json:"ticket_url"`
	Title     string     `json:"title"`
}

Added is one slice a pass added.

type Candidate

type Candidate struct {
	Source SourceName `json:"source"`
	// Key is the work's stable identity; the slice identifier derives from
	// it, so the same work is never added twice.
	Key   string `json:"key"`
	Title string `json:"title"`
	// TicketURL is the existing ticket the work delivers; empty when the
	// dreamer opens one.
	TicketURL string `json:"ticket_url,omitempty"`
	// Frequency is how many times a week the problem occurs, with how that
	// was measured.
	Frequency           float64 `json:"frequency"`
	FrequencyDerivation string  `json:"frequency_derivation"`
	// Fixed and Attempted are the measured outcomes of earlier work on the
	// same class: finished attempts and those that delivered.
	Fixed     int `json:"fixed"`
	Attempted int `json:"attempted"`
	// Known is true when earlier work on the class exists: exploration does
	// not weigh it.
	Known      bool     `json:"known"`
	Evidence   []string `json:"evidence"`
	Acceptance []string `json:"acceptance"`
}

Candidate is one piece of work a source proposes, with the inputs its rank is computed from and the typed acceptance and evidence its slice carries.

type Collector

type Collector struct {
	// contains filtered or unexported fields
}

Collector exports the dreamer's latest snapshot as the families above.

func (*Collector) Collect

func (exporter *Collector) Collect(metrics chan<- prometheus.Metric)

Collect sends the latest snapshot's samples.

func (*Collector) Describe

func (*Collector) Describe(descriptions chan<- *prometheus.Desc)

Describe sends every family's descriptor.

type Comparison

type Comparison struct {
	ComputedAt time.Time `json:"computed_at"`
	Week       string    `json:"week"`
	All        Partition `json:"all"`
	InSample   Partition `json:"in_sample"`
	HeldOut    Partition `json:"held_out"`
	// Explore is one plus the trailing complete weeks whose held-out rate
	// did not fall below the week before.
	Explore           float64 `json:"explore"`
	ExploreDerivation string  `json:"explore_derivation"`
}

Comparison is the held-out check: the weekly struggle factor on every run, on the in-sample runs the dreamer chooses work from and on the held-out runs it never reads, and the exploration weight the held-out series sets.

func Compare

func Compare(ctx context.Context, corpus stdfs.FS, location *time.Location, now time.Time) (Comparison, error)

Compare reads every run's event log in corpus, one run directory each, as the mining loop's measure reads it, and compares the partitions at now.

type Control

type Control struct {
	Action ControlAction `json:"action"`
	Reason string        `json:"reason"`
}

Control is one recorded control.

type ControlAction

type ControlAction string

ControlAction is a pause or a resume of the dreamer.

const (
	ControlPause  ControlAction = "pause"
	ControlResume ControlAction = "resume"
)

The controls.

type ControlInput

type ControlInput struct {
	Reason string `json:"reason" jsonschema:"why, recorded with the control"`
}

ControlInput is a pause or a resume of the dreamer.

type Decision

type Decision struct {
	Outcome  Outcome          `json:"outcome"`
	Reason   string           `json:"reason"`
	Target   Target           `json:"target"`
	Limits   []dispatch.Limit `json:"limits,omitempty"`
	Explore  float64          `json:"explore"`
	Ranked   []Ranked         `json:"ranked,omitempty"`
	Total    int              `json:"candidates"`
	Added    []Added          `json:"added,omitempty"`
	Released []string         `json:"released,omitempty"`
	// Findings are the sources that failed and anything skipped, each with
	// why.
	Findings []string `json:"findings,omitempty"`
}

Decision is one pass: what it read, why it did what it did, and what it added and released.

type DreamInput

type DreamInput struct{}

DreamInput asks for one pass now; it has no fields.

type Dreamer

type Dreamer struct {
	// contains filtered or unexported fields
}

Dreamer is the service. One goroutine owns its state and runs every pass, control and strategy in turn; the latest snapshot is published for reads.

func NewDreamer

func NewDreamer(options ...Option) (*Dreamer, error)

NewDreamer validates the whole option set before building the service. WithDispatcher and WithRecipe are required.

func (*Dreamer) Collector

func (dreamer *Dreamer) Collector() *Collector

Collector is the dreamer's exporter: register it on the binary's metrics registry.

func (*Dreamer) CurrentSnapshot

func (dreamer *Dreamer) CurrentSnapshot(_ context.Context, _ SnapshotInput) (Snapshot, error)

CurrentSnapshot reports the latest published snapshot.

func (*Dreamer) Dream

func (dreamer *Dreamer) Dream(ctx context.Context, _ DreamInput) (Decision, error)

Dream is one pass; it returns the decision it recorded.

func (*Dreamer) Pause

func (dreamer *Dreamer) Pause(ctx context.Context, input ControlInput) (Snapshot, error)

Pause stops every pass from adding or releasing anything until Resume.

func (*Dreamer) Register

func (dreamer *Dreamer) Register(router gin.IRouter)

Register mounts every operation's HTTP route on the caller's router.

func (*Dreamer) Resume

func (dreamer *Dreamer) Resume(ctx context.Context, input ControlInput) (Snapshot, error)

Resume lets the passes act again.

func (*Dreamer) Start

func (dreamer *Dreamer) Start(scope *runtime.Scope) error

Start replays the history and starts the owner goroutine on the scope.

func (*Dreamer) Strategize

func (dreamer *Dreamer) Strategize(ctx context.Context, _ StrategizeInput) (Comparison, error)

Strategize runs the held-out comparison and records the strategy it sets.

func (*Dreamer) Tools

func (dreamer *Dreamer) Tools() []csf.Option

Tools is every operation as an MCP tool for the CSF service.

func (*Dreamer) Triggers

func (dreamer *Dreamer) Triggers() []cronservice.Option

Triggers declare the pass and the daily strategy for the cron service.

type GitHubBacklog

type GitHubBacklog struct {
	// contains filtered or unexported fields
}

GitHubBacklog is CSF's own repository's issues over gh, run through the process capability.

func NewGitHubBacklog

func NewGitHubBacklog(launcher proc.ILauncher, repository string) (*GitHubBacklog, error)

NewGitHubBacklog grants one repository's issues; repository is its owner/name slug.

func (*GitHubBacklog) OpenIssue

func (backlog *GitHubBacklog) OpenIssue(ctx context.Context, title string, body string) (string, error)

OpenIssue opens one issue and returns its URL.

func (*GitHubBacklog) OpenIssues

func (backlog *GitHubBacklog) OpenIssues(ctx context.Context) ([]Issue, error)

OpenIssues lists the repository's open issues.

type IBacklog

type IBacklog interface {
	OpenIssue(ctx context.Context, title string, body string) (string, error)
}

IBacklog opens the ticket a generated slice delivers, in CSF's own repository. *GitHubBacklog satisfies it.

type IDispatcher

IDispatcher is what the dreamer needs of the slice dispatcher: its snapshot and graph, AddSlice, and Release for its own visibility holds. *dispatch.DispatchService satisfies it.

type Issue

type Issue struct {
	Number int64
	Title  string
	URL    string
	Body   string
	Labels []string
}

Issue is one open issue of CSF's own repository.

type Option

type Option func(dreamer *Dreamer) error

Option configures a Dreamer.

func WithBacklog

func WithBacklog(backlog IBacklog) Option

WithBacklog grants the ticket capability a generated slice without a ticket opens one through. Without it such a candidate is skipped.

func WithClock

func WithClock(clock harness.IClock) Option

WithClock replaces the clock.

func WithComparison

func WithComparison(compare func(ctx context.Context) (Comparison, error)) Option

WithComparison grants the held-out comparison Strategize runs.

func WithDispatcher

func WithDispatcher(dispatcher IDispatcher) Option

WithDispatcher grants the slice dispatcher the dreamer reads and adds through. Required.

func WithHistory

func WithHistory(records []Record) Option

WithHistory replays the records an earlier process wrote, oldest first: the pause, the strategy, the slices added and the tickets opened.

func WithLogger

func WithLogger(logger *slog.Logger) Option

WithLogger receives the service's records.

func WithRecipe

func WithRecipe(recipe *pb.AgentAssignmentRecipe) Option

WithRecipe is the recipe every generated slice starts from: its agent, model and workspace. Its task, when set, opens every brief. Required.

func WithRecordSink

func WithRecordSink(sink RecordSink) Option

WithRecordSink receives every record.

func WithSnapshotSink

func WithSnapshotSink(sink SnapshotSink) Option

WithSnapshotSink receives every published snapshot.

func WithSource

func WithSource(source Source) Option

WithSource registers one source of work.

type Outcome

type Outcome string

Outcome is how a pass ended.

const (
	OutcomePaused           Outcome = "paused"
	OutcomeDispatcherPaused Outcome = "dispatcher_paused"
	OutcomeLimited          Outcome = "limited"
	OutcomeFull             Outcome = "full"
	OutcomeFilled           Outcome = "filled"
	OutcomeDry              Outcome = "dry"
)

The outcomes of a pass.

type Outcomes

type Outcomes struct {
	Fixed     int `json:"fixed"`
	Attempted int `json:"attempted"`
}

Outcomes are finished attempts and the ones that delivered.

type Partition

type Partition struct {
	Name   string           `json:"name"`
	Runs   int              `json:"runs"`
	Weeks  []ouroboros.Rate `json:"weeks"`
	Factor ouroboros.Trend  `json:"factor"`
}

Partition is one side of the comparison: its runs, its weekly struggle rates and their compounding factor per week with the 95% interval.

type Ranked

type Ranked struct {
	Candidate
	Fixability           float64 `json:"fixability"`
	FixabilityDerivation string  `json:"fixability_derivation"`
	Explore              float64 `json:"explore"`
	Score                float64 `json:"score"`
	Rank                 int     `json:"rank"`
}

Ranked is one candidate with its score and every input it came from.

func Rank

func Rank(candidates []Candidate, outcomes map[SourceName]Outcomes, explore float64) []Ranked

Rank scores every candidate and orders them best first: score, then the source order the constants declare, then the key. outcomes are the dreamer's own slices' outcomes by source; explore weighs unknown classes.

type Record

type Record struct {
	At       time.Time   `json:"at"`
	Kind     RecordKind  `json:"kind"`
	Decision *Decision   `json:"decision,omitempty"`
	Control  *Control    `json:"control,omitempty"`
	Strategy *Comparison `json:"strategy,omitempty"`
}

Record is one typed record the dreamer writes: a pass's decision, a control, or a strategy with the comparison it followed.

type RecordKind

type RecordKind string

RecordKind says which part of a record is set.

const (
	RecordDecision RecordKind = "decision"
	RecordControl  RecordKind = "control"
	RecordStrategy RecordKind = "strategy"
)

The record kinds.

type RecordSink

type RecordSink func(record Record) error

RecordSink receives every record, on the dreamer's goroutine.

type Snapshot

type Snapshot struct {
	At          time.Time          `json:"at"`
	Paused      bool               `json:"paused"`
	PauseReason string             `json:"pause_reason,omitempty"`
	Explore     float64            `json:"explore"`
	Strategy    *Comparison        `json:"strategy,omitempty"`
	Generated   map[SourceName]int `json:"generated"`
	Decisions   []Record           `json:"decisions"`
}

Snapshot is what the Workbench shows of the dreamer: whether it is paused, the strategy and its comparison, how many slices each source generated, and the latest decisions, newest first.

type SnapshotInput

type SnapshotInput struct{}

SnapshotInput asks for the dreamer's snapshot; it has no fields.

type SnapshotSink

type SnapshotSink func(snapshot Snapshot) error

SnapshotSink receives every snapshot the dreamer publishes.

type Source

type Source struct {
	Name     string
	Fallback bool
	Find     func(ctx context.Context, world World) ([]Candidate, error)
}

Source is one registered origin of work. A fallback source is asked only when every primary source found nothing.

func Backlog

func Backlog(list func(ctx context.Context) ([]Issue, error), store ouroboros.IStore) Source

Backlog proposes the open issues of CSF's own repository that no slice delivers and no session claimed: mining tickets with no ready or running fixer are ungated findings, issues labelled ruling or intent are unenforced rulings, and the rest are unowned tickets. Language tickets are the operator's and are never taken. store, when given, supplies the mining tickets' fixers.

func HarnessGaps

func HarnessGaps(store ouroboros.IStore) Source

HarnessGaps is the fallback: the harness's own gaps measured from its records over the week. Refusals are the merge train's refused merges and the dispatcher's failed slices, grouped by their reason's first line; cost is the loop's fixer spend per miner that reached no ready pull request, highest first; the slowest are the miners whose fixers took longest on average.

func OntologySignals

func OntologySignals(files stdfs.FS, name string) Source

OntologySignals proposes every measured ontology alignment signal with a count above zero in the score record name, as tools/ontology-score.sh prints it, read from files. With no record there is nothing to propose.

func StruggleClasses

func StruggleClasses(store ouroboros.IStore) Source

StruggleClasses proposes the mining loop's struggle classes: the findings of the week on in-sample runs, grouped by miner and rule. Frequency is the class's count in the week; the earlier work is the loop's fixers for the miner. Held-out runs are never read.

type SourceName

type SourceName string

SourceName names where a candidate came from.

const (
	SourceStruggle SourceName = "struggle_class"
	SourceFinding  SourceName = "ungated_finding"
	SourceOntology SourceName = "ontology_signal"
	SourceTicket   SourceName = "unowned_ticket"
	SourceRuling   SourceName = "unenforced_ruling"
	SourceRefusal  SourceName = "harness_refusal"
	SourceCost     SourceName = "harness_cost"
	SourceSlowest  SourceName = "harness_slowest"
)

The sources, primary first, then the harness's own gaps.

type StrategizeInput

type StrategizeInput struct{}

StrategizeInput asks for the strategy to be recomputed now; it has no fields.

type Target

type Target struct {
	Ready    int `json:"ready"`
	Target   int `json:"target"`
	Capacity int `json:"capacity"`
	Running  int `json:"running"`
	// Completions is how many running sessions are expected to finish before
	// the next pass.
	Completions int `json:"completions"`
	// Measured is how many merged slices the median duration is over.
	Measured      int     `json:"measured"`
	MedianSeconds float64 `json:"median_seconds"`
	Derivation    string  `json:"derivation"`
}

Target is how long the ready frontier should be, with what it was derived from.

func DeriveTarget

func DeriveTarget(snapshot dispatch.Snapshot, durations []time.Duration, period time.Duration, visibility string) Target

DeriveTarget is the ready frontier the admission can use until the next pass: the launches it allows now (capacity less what runs) plus the running sessions expected to finish within one period, running × period / the median enqueue-to-merge duration of the merged slices. Before any slice has merged the second term is zero. ready counts the queued frontier and the dreamer's own slices held for visibility, which the next pass releases.

type World

type World struct {
	Now    time.Time
	Slices []*dispatchv1.SliceNode
}

World is what a source may read of the dispatcher at the pass that asks it: the time and every slice in the graph.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL