Documentation
¶
Index ¶
- Constants
- func IsConnectTimeout(err error) bool
- func PeerPID(ctx context.Context) int
- func Subscribe(socketPath string, params *SubscribeParams) (<-chan Event, func(), error)
- type AdmitPushParams
- type AdmitPushResult
- type CancelRunParams
- type CancelRunResult
- type Client
- func (c *Client) Call(method string, params interface{}, result interface{}) error
- func (c *Client) CallWithContext(ctx context.Context, method string, params interface{}, result interface{}, ...) error
- func (c *Client) CallWithTimeout(method string, params interface{}, result interface{}, timeout time.Duration) error
- func (c *Client) Close() error
- type ConnectTimeoutError
- type Event
- type EventClass
- type EventType
- type GateContextParams
- type GateContextResult
- type GetActiveRunParams
- type GetActiveRunResult
- type GetRunParams
- type GetRunResult
- type GetRunsForHeadParams
- type GetRunsParams
- type GetRunsResult
- type GetStepDiffParams
- type GetStepDiffResult
- type HandlerFunc
- type HealthParams
- type HealthResult
- type PushReceivedParams
- type PushReceivedResult
- type RPCError
- type Request
- type RerunParams
- type RerunResult
- type RespondParams
- type RespondResult
- type Response
- type RunInfo
- type Server
- func (s *Server) Close()
- func (s *Server) CloseListener()
- func (s *Server) Handle(method string, fn HandlerFunc)
- func (s *Server) HandleStream(method string, fn StreamHandlerFunc)
- func (s *Server) Listen(socketPath string) error
- func (s *Server) Serve(socketPath string) error
- func (s *Server) ServeReady() error
- type ShutdownParams
- type ShutdownResult
- type StepResultInfo
- type StreamFunc
- type StreamHandlerFunc
- type SubscribeParams
Constants ¶
const ( MethodPushReceived = "push_received" MethodGetRun = "get_run" MethodGetStepDiff = "get_step_diff" MethodGetRuns = "get_runs" MethodGetRunsForHead = "get_runs_for_head" MethodGetActiveRun = "get_active_run" MethodRerun = "rerun" MethodSubscribe = "subscribe" MethodRespond = "respond" MethodCancelRun = "cancel_run" MethodGateContext = "gate_context" MethodAdmitPush = "admit_push" MethodHealth = "health" MethodShutdown = "shutdown" )
JSON-RPC 2.0 method names.
const ( ErrParseError = -32700 ErrInvalidRequest = -32600 ErrMethodNotFound = -32601 ErrInvalidParams = -32602 ErrInternal = -32603 )
JSON-RPC 2.0 error codes.
const ( // DefaultDialTimeout is the read deadline used for the daemon health check // dial made by callers outside this package (see internal/daemon/selfexec.go). DefaultDialTimeout = 250 * time.Millisecond )
Variables ¶
This section is empty.
Functions ¶
func IsConnectTimeout ¶
IsConnectTimeout reports whether err was caused by a bounded IPC connect timeout.
func PeerPID ¶
PeerPID returns the OS-authenticated process ID of the local IPC client. Zero means the transport cannot authenticate a peer process on this platform.
func Subscribe ¶
func Subscribe(socketPath string, params *SubscribeParams) (<-chan Event, func(), error)
Subscribe opens a dedicated connection and subscribes to events for a run. Returns an event channel, a cancel function (to stop and clean up), and an error. The channel is closed when the run completes, the connection drops, or cancel is called.
Types ¶
type AdmitPushParams ¶
type AdmitPushParams struct {
Gate string `json:"gate"`
}
AdmitPushParams asks whether a local receive hook's authenticated process ancestry is allowed to mutate a managed gate ref.
type AdmitPushResult ¶
type AdmitPushResult struct {
Context GateContextResult `json:"context"`
}
AdmitPushResult is returned before a receive hook permits ref mutation.
type CancelRunParams ¶
type CancelRunParams struct {
RunID string `json:"run_id"`
}
CancelRunParams cancels an active pipeline run.
type CancelRunResult ¶
type CancelRunResult struct {
OK bool `json:"ok"`
}
CancelRunResult confirms the run cancellation request was accepted.
type Client ¶
type Client struct {
// contains filtered or unexported fields
}
Client connects to the IPC server over the platform transport.
func (*Client) Call ¶
Call sends a JSON-RPC request and waits for the response. The result is unmarshaled into the provided pointer. If the server returns a JSON-RPC error, it is returned as *RPCError.
func (*Client) CallWithContext ¶
func (c *Client) CallWithContext(ctx context.Context, method string, params interface{}, result interface{}, timeout time.Duration) error
CallWithContext is CallWithTimeout with cancellation support.
type ConnectTimeoutError ¶
ConnectTimeoutError reports a daemon IPC connect attempt that exceeded the bounded client timeout.
func (*ConnectTimeoutError) Error ¶
func (e *ConnectTimeoutError) Error() string
func (*ConnectTimeoutError) Timeout ¶
func (e *ConnectTimeoutError) Timeout() bool
func (*ConnectTimeoutError) Unwrap ¶
func (e *ConnectTimeoutError) Unwrap() error
type Event ¶
type Event struct {
Type EventType `json:"type"`
RunID string `json:"run_id"`
RepoID string `json:"repo_id"`
StepName *types.StepName `json:"step_name,omitempty"`
Status *string `json:"status,omitempty"`
Error *string `json:"error,omitempty"`
Stream *string `json:"stream,omitempty"`
Content *string `json:"content,omitempty"`
Branch *string `json:"branch,omitempty"`
Findings *string `json:"findings,omitempty"` // JSON-encoded findings for step_completed events
ReportedFindings *int `json:"reported_findings,omitempty"`
FixedFindings *int `json:"fixed_findings,omitempty"`
DurationMS *int64 `json:"duration_ms,omitempty"` // execution-only duration for step events
PRURL *string `json:"pr_url,omitempty"` // PR URL for run_updated/run_completed events
// StateRev is the daemon-assigned monotonic revision of the run state
// this event reflects, or zero for activity. A consumer applies a state
// delta only when StateRev exceeds the revision it has already applied,
// which makes a delta queued before an authoritative snapshot an
// idempotent no-op after it.
StateRev int64 `json:"state_rev,omitempty"`
CIReady *bool `json:"ci_ready,omitempty"`
CIReadyNoCI *bool `json:"ci_ready_no_ci,omitempty"`
}
Event is a real-time update sent to subscribers.
type EventClass ¶
type EventClass int
EventClass is the loss-tolerance taxonomy for the run event stream. It is the single owner of "may this event be dropped?", so the daemon's overflow policy and every consumer's reconciliation policy agree by construction rather than by two lists of event names drifting apart.
const ( // ClassActivity is ephemeral output with no state effect. Losing it // cannot make a consumer render a state the daemon does not hold, so it // is the only class that may be dropped under pressure. ClassActivity EventClass = iota // ClassState is a state transition. Its payload is a wakeup hint only - // every field is reconstructable from get_run - but the *fact* that // state changed must reach every subscriber, or the consumer keeps // rendering state the daemon has already left behind. ClassState // ClassControl is stream metadata generated by the broker itself, never // by the executor. It is delivered ahead of queued payload. ClassControl )
func ClassOf ¶
func ClassOf(t EventType) EventClass
ClassOf classifies an event type for overflow and reconciliation handling.
Unknown types fail safe to ClassState: an event a future producer adds and this build does not recognise is never silently dropped. Only the types explicitly named here are droppable.
type EventType ¶
type EventType string
EventType identifies the kind of event.
const ( EventRunCreated EventType = "run_created" EventRunUpdated EventType = "run_updated" EventRunCompleted EventType = "run_completed" EventCIReadinessChanged EventType = "ci_readiness_changed" EventStepStarted EventType = "step_started" EventStepCompleted EventType = "step_completed" EventStepsReset EventType = "steps_reset" EventLogChunk EventType = "log_chunk" // EventStreamGap tells a subscriber that the daemon coalesced at least // one state transition away under buffer pressure. StateRev is the // highest revision folded into it. The subscriber must read authoritative // state once; the frame carries no payload of its own. EventStreamGap EventType = "stream_gap" )
type GateContextParams ¶
type GateContextParams struct {
CWD string `json:"cwd,omitempty"`
MarkerPresent bool `json:"marker_present,omitempty"`
}
GateContextParams asks the daemon to classify the authenticated caller. CWD and MarkerPresent are evidence only; peer PID comes from the transport.
type GateContextResult ¶
type GateContextResult struct {
Nested bool `json:"nested"`
ManagedGit bool `json:"managed_git,omitempty"`
AgentDescendant bool `json:"agent_descendant,omitempty"`
DaemonDescendant bool `json:"daemon_descendant,omitempty"`
MarkerPresent bool `json:"marker_present,omitempty"`
RunID string `json:"run_id,omitempty"`
Phase types.StepName `json:"phase,omitempty"`
}
GateContextResult is the privacy-safe execution-context classification.
type GetActiveRunParams ¶
type GetActiveRunParams struct {
RepoID string `json:"repo_id"`
Branch string `json:"branch,omitempty"`
}
GetActiveRunParams requests the active run for a repo. When Branch is set, runs on that branch are preferred.
type GetActiveRunResult ¶
type GetActiveRunResult struct {
Run *RunInfo `json:"run,omitempty"`
}
GetActiveRunResult wraps the active run (nil if none).
type GetRunParams ¶
type GetRunParams struct {
RunID string `json:"run_id"`
}
GetRunParams requests a single run by ID.
type GetRunResult ¶
type GetRunResult struct {
Run *RunInfo `json:"run"`
}
GetRunResult wraps a single run.
type GetRunsForHeadParams ¶
type GetRunsForHeadParams struct {
RepoID string `json:"repo_id"`
Branch string `json:"branch"`
HeadSHA string `json:"head_sha"`
}
GetRunsForHeadParams requests the runs for a repo on an exact branch and head SHA. It backs a lightweight lookup that avoids scanning the repo's whole run history, so a caller polling for the run created by a specific push does not re-fetch every run (and its steps) on each poll.
type GetRunsParams ¶
type GetRunsParams struct {
RepoID string `json:"repo_id"`
}
GetRunsParams requests all runs for a repo.
type GetRunsResult ¶
type GetRunsResult struct {
Runs []RunInfo `json:"runs"`
}
GetRunsResult wraps a list of runs.
type GetStepDiffParams ¶
type GetStepDiffParams struct {
RunID string `json:"run_id"`
}
GetStepDiffParams requests the working-tree diff for a run parked at a fix-review gate. The diff is derived on demand from the run's worktree and is never stored, so it is the reconstruction authority for the one piece of gate context that is not persisted.
type GetStepDiffResult ¶
type GetStepDiffResult struct {
Diff string `json:"diff"`
Truncated bool `json:"truncated,omitempty"`
}
GetStepDiffResult carries a bounded working-tree diff. Truncated reports that the diff exceeded the response budget and was cut, so a very large change degrades to a partial view instead of an oversized frame.
type HandlerFunc ¶
type HandlerFunc func(ctx context.Context, params json.RawMessage) (interface{}, error)
HandlerFunc processes a JSON-RPC request and returns a result or error.
type HealthParams ¶
type HealthParams struct{}
HealthParams has no fields but exists for consistency.
type HealthResult ¶
type HealthResult struct {
Status string `json:"status"`
}
HealthResult confirms the daemon is alive.
type PushReceivedParams ¶
type PushReceivedParams struct {
// Gate is the absolute path to the gate bare repo.
Gate string `json:"gate"`
Ref string `json:"ref"`
Old string `json:"old"`
New string `json:"new"`
SkipSteps []types.StepName `json:"skip_steps,omitempty"`
Intent string `json:"intent,omitempty"`
}
PushReceivedParams are sent by the post-receive hook when a push arrives.
Intent, when set, is an agent-supplied description of the change. It is stamped onto the run so the intent step uses it verbatim instead of inferring intent from local transcripts.
type PushReceivedResult ¶
type PushReceivedResult struct {
RunID string `json:"run_id"`
}
PushReceivedResult confirms the push was accepted.
type Request ¶
type Request struct {
JSONRPC string `json:"jsonrpc"`
Method string `json:"method"`
Params json.RawMessage `json:"params,omitempty"`
ID int64 `json:"id"`
}
Request is a JSON-RPC 2.0 request.
func NewRequest ¶
NewRequest creates a JSON-RPC 2.0 request with an auto-incremented ID.
type RerunParams ¶
type RerunParams struct {
RepoID string `json:"repo_id"`
Branch string `json:"branch"`
PreviousRunID string `json:"previous_run_id,omitempty"`
SkipSteps []types.StepName `json:"skip_steps,omitempty"`
Intent string `json:"intent,omitempty"`
}
RerunParams requests a new run for the latest gate head on a branch. Intent, when set, overrides inherited intent and fresh inference. When empty, the daemon inherits authoritative intent from the selected prior run or leaves the new run to perform fresh inference.
type RerunResult ¶
type RerunResult struct {
RunID string `json:"run_id"`
}
RerunResult confirms a rerun was created.
type RespondParams ¶
type RespondParams struct {
RunID string `json:"run_id"`
Step types.StepName `json:"step"`
Action types.ApprovalAction `json:"action"`
FindingIDs []string `json:"finding_ids,omitempty"`
Instructions map[string]string `json:"instructions,omitempty"`
AddedFindings []types.Finding `json:"added_findings,omitempty"`
}
RespondParams sends a user action for a step awaiting approval.
Instructions carries optional per-finding notes keyed by finding ID, which the daemon attaches to the corresponding finding before dispatching a fix. AddedFindings carries user-authored findings that are merged into the round alongside agent-produced ones. Both fields only apply when Action triggers a fix round.
type RespondResult ¶
type RespondResult struct {
OK bool `json:"ok"`
}
RespondResult confirms the action was accepted.
type Response ¶
type Response struct {
JSONRPC string `json:"jsonrpc"`
Result json.RawMessage `json:"result,omitempty"`
Error *RPCError `json:"error,omitempty"`
ID int64 `json:"id"`
}
Response is a JSON-RPC 2.0 response.
func NewErrorResponse ¶
NewErrorResponse creates an error JSON-RPC 2.0 response.
func NewResponse ¶
NewResponse creates a successful JSON-RPC 2.0 response.
type RunInfo ¶
type RunInfo struct {
ID string `json:"id"`
RepoID string `json:"repo_id"`
Branch string `json:"branch"`
HeadSHA string `json:"head_sha"`
SubmittedHeadSHA *string `json:"submitted_head_sha,omitempty"`
BaseSHA string `json:"base_sha"`
Status types.RunStatus `json:"status"`
PRURL *string `json:"pr_url,omitempty"`
Error *string `json:"error,omitempty"`
CIReady bool `json:"ci_ready,omitempty"`
CIReadyNoCI bool `json:"ci_ready_no_ci,omitempty"`
// AwaitingAgent is true while the run is parked at a gate awaiting the
// driving agent's response. AwaitingAgentSince is the unix-seconds time it
// parked, so a supervisor can read "parked for N seconds" in one call. Both
// are observability only and clear the moment the agent responds.
AwaitingAgent bool `json:"awaiting_agent,omitempty"`
AwaitingAgentSince *int64 `json:"awaiting_agent_since,omitempty"`
Steps []StepResultInfo `json:"steps,omitempty"`
// StateRev is the monotonic run-state revision this snapshot is at least
// as new as. It is sampled before the database read, so every event at or
// below it is already reflected here and every event above it still
// applies on top.
StateRev int64 `json:"state_rev,omitempty"`
CreatedAt int64 `json:"created_at"`
UpdatedAt int64 `json:"updated_at"`
}
RunInfo is the IPC representation of a pipeline run.
type Server ¶
type Server struct {
// contains filtered or unexported fields
}
Server listens on an IPC endpoint and dispatches JSON-RPC requests.
func (*Server) CloseListener ¶
func (s *Server) CloseListener()
CloseListener closes the underlying listener without signaling server shutdown. This causes Accept to return net.ErrClosed, which the server detects and exits cleanly.
func (*Server) Handle ¶
func (s *Server) Handle(method string, fn HandlerFunc)
Handle registers a handler for a JSON-RPC method.
func (*Server) HandleStream ¶
func (s *Server) HandleStream(method string, fn StreamHandlerFunc)
HandleStream registers a streaming handler for a JSON-RPC method. When this method is called, the server sends an initial OK response, then hands the connection to the handler for streaming. The connection closes when the handler returns.
func (*Server) Listen ¶
Listen binds the IPC endpoint without starting the accept loop. It lets the daemon measure bind separately, then start serving and prove real health before announcing readiness.
func (*Server) Serve ¶
Serve starts listening on the given IPC endpoint path. It blocks until Close is called, then returns nil.
func (*Server) ServeReady ¶
ServeReady accepts requests from an endpoint already bound by Listen.
type ShutdownParams ¶
type ShutdownParams struct{}
ShutdownParams has no fields but exists for consistency.
type ShutdownResult ¶
type ShutdownResult struct {
OK bool `json:"ok"`
}
ShutdownResult confirms shutdown was initiated.
type StepResultInfo ¶
type StepResultInfo struct {
ID string `json:"id"`
RunID string `json:"run_id"`
StepName types.StepName `json:"step_name"`
StepOrder int `json:"step_order"`
Status types.StepStatus `json:"status"`
ExitCode *int `json:"exit_code,omitempty"`
DurationMS *int64 `json:"duration_ms,omitempty"`
FindingsJSON *string `json:"findings_json,omitempty"`
ReportedFindings int `json:"reported_findings,omitempty"`
FixedFindings int `json:"fixed_findings,omitempty"`
// FixSummaries holds one entry per fix round the pipeline ran for this
// step, in round order: the agent's one-line fix summary, or "" when the
// round recorded none. Agent surfaces use it to report applied fixes.
FixSummaries []string `json:"fix_summaries,omitempty"`
RoundCount int `json:"round_count,omitempty"`
FixRoundCount int `json:"fix_round_count,omitempty"`
AutoFixLimit int `json:"auto_fix_limit,omitempty"`
PendingFixSource string `json:"pending_fix_source,omitempty"`
// ConvergenceJSON is the review step's persisted convergence report
// (internal/convergence.Report), so gate renders reached over IPC carry
// the per-round history without a separate database read.
ConvergenceJSON *string `json:"convergence_json,omitempty"`
Error *string `json:"error,omitempty"`
StartedAt *int64 `json:"started_at,omitempty"`
CompletedAt *int64 `json:"completed_at,omitempty"`
LastActivityAt *int64 `json:"last_activity_at,omitempty"`
LastActivity *string `json:"last_activity,omitempty"`
AgentPID *int `json:"agent_pid,omitempty"`
}
StepResultInfo is the IPC representation of a step result.
type StreamFunc ¶
StreamFunc owns a prepared streaming connection. send writes a JSON object; the function should block until streaming is complete and return when send reports a disconnected client.
type StreamHandlerFunc ¶
type StreamHandlerFunc func(ctx context.Context, params json.RawMessage) (StreamFunc, error)
StreamHandlerFunc prepares a stream before the server acknowledges the subscription. This closes the subscribe-then-reconcile race: callers cannot perform their first reconciliation until the handler has registered its event source. The returned StreamFunc runs after the acknowledgement; preparation resources must also be released when ctx is cancelled because an acknowledgement write can fail before StreamFunc starts.
type SubscribeParams ¶
type SubscribeParams struct {
RunID string `json:"run_id"`
}
SubscribeParams starts an event stream for a run.