Documentation
¶
Overview ¶
Package control provides the adapters that plug Fort's deterministic components into the control-plane ports (ui.Dispatcher, ui.FlowRunner).
This is the composition seam for "control plane, optionally with execution":
- EngineDispatcher / FlowExecutor wrap the real router + DAG engine (full mode).
- QueueDispatcher boards tasks with no execution plane at all (control-only mode).
cmd/fort picks which adapters to wire; the ui module only ever sees the ports.
Index ¶
- type CapabilityCoordinator
- func (c *CapabilityCoordinator) Capabilities() (corecap.Snapshot, uint64)
- func (c *CapabilityCoordinator) Current() (corecap.Snapshot, uint64)
- func (c *CapabilityCoordinator) Refresh(ctx context.Context, mode corecap.RefreshMode, adapters []string) (corecap.Snapshot, uint64, error)
- func (c *CapabilityCoordinator) RefreshMachine(ctx context.Context, target string, mode corecap.RefreshMode, ...) (corecap.MachineInventory, error)
- type CapabilityCoordinatorOptions
- type CapabilitySnapshotSource
- type ConversationActivity
- type ConversationSeatSource
- type ConversationService
- func (s *ConversationService) AddConversationParticipant(_ context.Context, conversationID, seatID string) (conversation.Participant, error)
- func (s *ConversationService) CancelTarget(_ context.Context, targetID string) error
- func (s *ConversationService) Close()
- func (s *ConversationService) ConversationSeats(context.Context) ([]conversation.Seat, error)
- func (s *ConversationService) ConversationTargetActive(id string) bool
- func (s *ConversationService) CreateConversation(_ context.Context, projectID, title string, seatIDs []string) (store.ConversationDetail, error)
- func (s *ConversationService) CreateProject(_ context.Context, name string) (conversation.Project, error)
- func (s *ConversationService) DeleteConversation(_ context.Context, id string) error
- func (s *ConversationService) DeleteProject(_ context.Context, id string) error
- func (s *ConversationService) GetConversation(_ context.Context, id string) (store.ConversationDetail, error)
- func (s *ConversationService) ListConversations(_ context.Context, scope string) ([]conversation.Conversation, error)
- func (s *ConversationService) ListProjects(context.Context) ([]conversation.Project, error)
- func (s *ConversationService) MoveConversation(_ context.Context, id, projectID string) error
- func (s *ConversationService) PostTurn(_ context.Context, conversationID, clientTurnID, text string, ...) (conversation.TurnResult, error)
- func (s *ConversationService) RemoveConversationParticipant(_ context.Context, conversationID, participantID string) error
- func (s *ConversationService) RenameConversation(_ context.Context, id, title string) error
- func (s *ConversationService) RenameProject(_ context.Context, id, name string) error
- func (s *ConversationService) RetryTarget(_ context.Context, targetID string) (conversation.Target, error)
- func (s *ConversationService) SetConversationState(_ context.Context, id string, state conversation.ConversationState) error
- func (s *ConversationService) Wait()
- type CreateConversationRequest
- type EngineDispatcher
- type FlowExecutor
- func (f FlowExecutor) Approve(runID, nodeID, edit string) error
- func (f FlowExecutor) Plan(flowID string) []ui.FlowNode
- func (f FlowExecutor) Reject(runID, nodeID, note string) error
- func (f FlowExecutor) ResumeFlow(ctx context.Context, flowID, runID string) (ui.RunResult, error)
- func (f FlowExecutor) ResumeFlowAsync(ctx context.Context, flowID, runID string) error
- func (f FlowExecutor) StartFlow(ctx context.Context, flowID, runID, payload string) (ui.RunResult, error)
- func (f FlowExecutor) StartFlowAsync(ctx context.Context, flowID, runID, payload string) (ui.RunResult, error)
- func (f FlowExecutor) StartPlaybook(ctx context.Context, route ui.RoutePreview, runID, direction string) (ui.PlaybookRunResult, error)
- func (f FlowExecutor) StartPlaybookAsync(ctx context.Context, route ui.RoutePreview, runID, direction string) (ui.PlaybookRunResult, error)
- func (f FlowExecutor) WithPlaybooks(catalog *PlaybookCatalog) FlowExecutor
- type LocalCapabilityRegistry
- type PeerCapabilityClient
- type Planner
- type PlaybookCatalog
- func (c *PlaybookCatalog) Duplicate(ctx context.Context, id string) (ui.Playbook, error)
- func (c *PlaybookCatalog) List(ctx context.Context) ([]ui.Playbook, error)
- func (c *PlaybookCatalog) Route(ctx context.Context, request ui.RouteRequest) (ui.RoutePreview, error)
- func (c *PlaybookCatalog) Save(ctx context.Context, definition ui.Playbook) (ui.Playbook, error)
- type PostConversationTurnRequest
- type QueueDispatcher
- type Roster
- type ScheduleCreator
- type ScheduleService
- type SnapshotConversationSeats
- type TodayScheduleSource
- type TodayService
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type CapabilityCoordinator ¶ added in v0.13.0
type CapabilityCoordinator struct {
// contains filtered or unexported fields
}
CapabilityCoordinator refreshes local and peer inventories concurrently, binds every peer payload to registry-owned identity/rank, and publishes one normalized snapshot generation atomically.
func NewCapabilityCoordinator ¶ added in v0.13.0
func NewCapabilityCoordinator(options CapabilityCoordinatorOptions) (*CapabilityCoordinator, error)
func (*CapabilityCoordinator) Capabilities ¶ added in v0.13.0
func (c *CapabilityCoordinator) Capabilities() (corecap.Snapshot, uint64)
Capabilities implements ui.CapabilityLister without exposing refresh or private probe controls to the presentation layer.
func (*CapabilityCoordinator) Current ¶ added in v0.13.0
func (c *CapabilityCoordinator) Current() (corecap.Snapshot, uint64)
func (*CapabilityCoordinator) Refresh ¶ added in v0.13.0
func (c *CapabilityCoordinator) Refresh(ctx context.Context, mode corecap.RefreshMode, adapters []string) (corecap.Snapshot, uint64, error)
func (*CapabilityCoordinator) RefreshMachine ¶ added in v0.13.0
func (c *CapabilityCoordinator) RefreshMachine(ctx context.Context, target string, mode corecap.RefreshMode, adapters []string) (corecap.MachineInventory, error)
RefreshMachine refreshes only the already-selected target. Dispatch uses this path so an unrelated slow or unavailable peer cannot delay provider startup or influence deterministic placement.
type CapabilityCoordinatorOptions ¶ added in v0.13.0
type CapabilityCoordinatorOptions struct {
Live *machines.Live
LocalName string
Local LocalCapabilityRegistry
Peers PeerCapabilityClient
Now func() time.Time
}
type CapabilitySnapshotSource ¶ added in v0.13.0
type ConversationActivity ¶ added in v0.13.0
type ConversationSeatSource ¶ added in v0.13.0
type ConversationSeatSource interface {
ConversationSeats() []conversation.Seat
}
func FakeConversationSeats ¶ added in v0.13.0
func FakeConversationSeats(machine string) ConversationSeatSource
FakeConversationSeats is test/demo plumbing for FORT_FAKE only. Production composition must use SnapshotConversationSeats so static claims never become ready seats.
type ConversationService ¶ added in v0.13.0
type ConversationService struct {
// contains filtered or unexported fields
}
ConversationService is the durable coordinator for projects, shared conversations, target fan-out, attribution, retry, and cancellation.
func NewConversationService ¶ added in v0.13.0
func NewConversationService(st *store.Store, rt runtime.Runtime, seats ConversationSeatSource, workdir string) *ConversationService
func (*ConversationService) AddConversationParticipant ¶ added in v0.13.0
func (s *ConversationService) AddConversationParticipant(_ context.Context, conversationID, seatID string) (conversation.Participant, error)
func (*ConversationService) CancelTarget ¶ added in v0.13.0
func (s *ConversationService) CancelTarget(_ context.Context, targetID string) error
func (*ConversationService) Close ¶ added in v0.13.0
func (s *ConversationService) Close()
Close stops active providers and joins their persistence loops.
func (*ConversationService) ConversationSeats ¶ added in v0.13.0
func (s *ConversationService) ConversationSeats(context.Context) ([]conversation.Seat, error)
func (*ConversationService) ConversationTargetActive ¶ added in v0.13.0
func (s *ConversationService) ConversationTargetActive(id string) bool
func (*ConversationService) CreateConversation ¶ added in v0.13.0
func (s *ConversationService) CreateConversation(_ context.Context, projectID, title string, seatIDs []string) (store.ConversationDetail, error)
func (*ConversationService) CreateProject ¶ added in v0.13.0
func (s *ConversationService) CreateProject(_ context.Context, name string) (conversation.Project, error)
func (*ConversationService) DeleteConversation ¶ added in v0.13.0
func (s *ConversationService) DeleteConversation(_ context.Context, id string) error
func (*ConversationService) DeleteProject ¶ added in v0.13.0
func (s *ConversationService) DeleteProject(_ context.Context, id string) error
func (*ConversationService) GetConversation ¶ added in v0.13.0
func (s *ConversationService) GetConversation(_ context.Context, id string) (store.ConversationDetail, error)
func (*ConversationService) ListConversations ¶ added in v0.13.0
func (s *ConversationService) ListConversations(_ context.Context, scope string) ([]conversation.Conversation, error)
func (*ConversationService) ListProjects ¶ added in v0.13.0
func (s *ConversationService) ListProjects(context.Context) ([]conversation.Project, error)
func (*ConversationService) MoveConversation ¶ added in v0.13.0
func (s *ConversationService) MoveConversation(_ context.Context, id, projectID string) error
func (*ConversationService) PostTurn ¶ added in v0.13.0
func (s *ConversationService) PostTurn(_ context.Context, conversationID, clientTurnID, text string, targetParticipantIDs []string) (conversation.TurnResult, error)
func (*ConversationService) RemoveConversationParticipant ¶ added in v0.13.0
func (s *ConversationService) RemoveConversationParticipant(_ context.Context, conversationID, participantID string) error
func (*ConversationService) RenameConversation ¶ added in v0.13.0
func (s *ConversationService) RenameConversation(_ context.Context, id, title string) error
func (*ConversationService) RenameProject ¶ added in v0.13.0
func (s *ConversationService) RenameProject(_ context.Context, id, name string) error
func (*ConversationService) RetryTarget ¶ added in v0.13.0
func (s *ConversationService) RetryTarget(_ context.Context, targetID string) (conversation.Target, error)
func (*ConversationService) SetConversationState ¶ added in v0.13.0
func (s *ConversationService) SetConversationState(_ context.Context, id string, state conversation.ConversationState) error
func (*ConversationService) Wait ¶ added in v0.13.0
func (s *ConversationService) Wait()
Wait joins every target accepted before the call. Shutdown and tests use it before closing the durable store that target completions persist to.
type CreateConversationRequest ¶ added in v0.13.0
type EngineDispatcher ¶
type EngineDispatcher struct {
// contains filtered or unexported fields
}
EngineDispatcher routes + dispatches via the deterministic engine.
func NewEngineDispatcher ¶
func NewEngineDispatcher(e *engine.Engine) EngineDispatcher
NewEngineDispatcher adapts an engine to ui.Dispatcher.
type FlowExecutor ¶
type FlowExecutor struct {
// contains filtered or unexported fields
}
FlowExecutor adapts the DAG executor to ui.FlowRunner (id-based).
func NewFlowExecutor ¶
func NewFlowExecutor(x *graph.Executor, flows []graph.Flow) FlowExecutor
NewFlowExecutor adapts a graph executor + flow set to ui.FlowRunner.
func (FlowExecutor) Approve ¶
func (f FlowExecutor) Approve(runID, nodeID, edit string) error
Approve records a gate approval.
func (FlowExecutor) Plan ¶ added in v0.11.0
func (f FlowExecutor) Plan(flowID string) []ui.FlowNode
Plan exposes a flow's node list to the control plane (spec 033).
func (FlowExecutor) Reject ¶
func (f FlowExecutor) Reject(runID, nodeID, note string) error
Reject records a gate rejection with an optional redirect note.
func (FlowExecutor) ResumeFlow ¶
ResumeFlow resumes the named flow.
func (FlowExecutor) ResumeFlowAsync ¶ added in v0.13.0
func (f FlowExecutor) ResumeFlowAsync(ctx context.Context, flowID, runID string) error
ResumeFlowAsync validates the immutable flow before scheduling one detached continuation of the durable run.
func (FlowExecutor) StartFlow ¶
func (f FlowExecutor) StartFlow(ctx context.Context, flowID, runID, payload string) (ui.RunResult, error)
StartFlow starts the named flow.
func (FlowExecutor) StartFlowAsync ¶ added in v0.13.0
func (f FlowExecutor) StartFlowAsync(ctx context.Context, flowID, runID, payload string) (ui.RunResult, error)
StartFlowAsync persists the run, then walks the flow once in the background.
func (FlowExecutor) StartPlaybook ¶ added in v0.12.0
func (f FlowExecutor) StartPlaybook(ctx context.Context, route ui.RoutePreview, runID, direction string) (ui.PlaybookRunResult, error)
StartPlaybook compiles and starts an immutable route preview. The canonical flow id carries every field needed to reconstruct the exact flow after a process restart, while its delivery suffix lets operational views exclude answer-only history.
func (FlowExecutor) StartPlaybookAsync ¶ added in v0.13.0
func (f FlowExecutor) StartPlaybookAsync(ctx context.Context, route ui.RoutePreview, runID, direction string) (ui.PlaybookRunResult, error)
StartPlaybookAsync persists the exact canonical playbook flow, then executes it once in the background. Quick-answer output is delivered later through the existing run detail and event stream rather than holding the HTTP response.
func (FlowExecutor) WithPlaybooks ¶ added in v0.12.0
func (f FlowExecutor) WithPlaybooks(catalog *PlaybookCatalog) FlowExecutor
WithPlaybooks enables immutable dynamic playbook flows on the same executor value used for static ui.FlowRunner operations.
type LocalCapabilityRegistry ¶ added in v0.13.0
type LocalCapabilityRegistry interface {
Current() corecap.NodeInventory
Refresh(context.Context, corecap.RecheckRequest) (corecap.NodeInventory, error)
}
type PeerCapabilityClient ¶ added in v0.13.0
type PeerCapabilityClient interface {
Refresh(context.Context, string, string, corecap.RecheckRequest) (corecap.NodeInventory, error)
}
type Planner ¶ added in v0.8.0
type Planner struct {
// contains filtered or unexported fields
}
Planner implements ui.Planner: it runs a planner agent to decompose a goal into backlog sub-tasks (spec 026). The planner is a normal, visible run; Breakdown returns its id immediately and the sub-tasks are created when the run completes.
func NewPlanner ¶ added in v0.8.0
NewPlanner adapts the engine + store to ui.Planner. defaultAgent is the planner agent used when a breakdown request doesn't specify one (FORT_PLANNER; falls back to "claude").
func (Planner) Breakdown ¶ added in v0.8.0
Breakdown dispatches the planner run and returns its id; a goroutine ingests its output into the backlog once it completes. NOTE (spec 026 crash window): if the process stops between the run finishing and ingest writing items, the sub-tasks are not written — re-run the breakdown.
type PlaybookCatalog ¶ added in v0.12.0
type PlaybookCatalog struct {
// contains filtered or unexported fields
}
PlaybookCatalog adapts immutable store revisions to the UI playbook port. Definitions remain encoded as their UI wire form so display-only fields such as stage descriptions survive validation and round trips unchanged.
func NewPlaybookCatalog ¶ added in v0.12.0
func NewPlaybookCatalog(st *store.Store) *PlaybookCatalog
NewPlaybookCatalog builds the durable playbook catalog and atomically seeds it before any route preview can be served. Routing remains read-only.
func (*PlaybookCatalog) Duplicate ¶ added in v0.12.0
Duplicate appends a disabled, non-default copy with a fresh deterministic slug. A duplicate cannot change automatic routing until explicitly enabled.
func (*PlaybookCatalog) List ¶ added in v0.12.0
List returns the latest immutable revision for every playbook.
func (*PlaybookCatalog) Route ¶ added in v0.12.0
func (c *PlaybookCatalog) Route(ctx context.Context, request ui.RouteRequest) (ui.RoutePreview, error)
Route resolves a pure deterministic preview. An explicit revision loads that exact immutable row; the optional plan gate override only changes the returned execution snapshot and never persists a definition.
type PostConversationTurnRequest ¶ added in v0.13.0
type QueueDispatcher ¶
type QueueDispatcher struct {
// contains filtered or unexported fields
}
QueueDispatcher boards a task as a "queued" run with no execution plane.
func NewQueueDispatcher ¶
func NewQueueDispatcher(s *store.Store) QueueDispatcher
NewQueueDispatcher adapts a store to ui.Dispatcher for control-only mode.
type Roster ¶
type Roster struct {
// contains filtered or unexported fields
}
Roster adapts a live machine registry pointer to ui.MachineLister and tracks peer reachability by polling each machine's /health (spec 022). It reads the Live pointer on every call, so a registry installed after construction (e.g. a mesh enrollment) is visible without a restart (spec 024). The local machine is reachable by definition; peers start unreachable until first probed.
func (*Roster) Machines ¶
func (r *Roster) Machines() []ui.MachineStatus
Machines implements ui.MachineLister.
type ScheduleCreator ¶ added in v0.13.0
type ScheduleCreator interface {
Create(context.Context, scheduler.Definition) error
}
type ScheduleService ¶ added in v0.13.0
type ScheduleService struct {
// contains filtered or unexported fields
}
func NewScheduleService ¶ added in v0.13.0
func NewScheduleService(next ScheduleCreator, flowIDs []string) *ScheduleService
func (*ScheduleService) Create ¶ added in v0.13.0
func (s *ScheduleService) Create(ctx context.Context, definition scheduler.Definition) error
type SnapshotConversationSeats ¶ added in v0.13.0
type SnapshotConversationSeats struct {
Source CapabilitySnapshotSource
}
SnapshotConversationSeats projects the verified capability inventory into exact profile + provider/model + machine seats.
func (SnapshotConversationSeats) ConversationSeats ¶ added in v0.13.0
func (s SnapshotConversationSeats) ConversationSeats() []conversation.Seat
type TodayScheduleSource ¶ added in v0.13.0
type TodayService ¶ added in v0.13.0
type TodayService struct {
// contains filtered or unexported fields
}
func NewTodayService ¶ added in v0.13.0
func NewTodayService(st *store.Store, schedules TodayScheduleSource, activity ConversationActivity) *TodayService