ipc

package
v1.55.0 Latest Latest
Warning

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

Go to latest
Published: Sep 11, 2026 License: MIT Imports: 17 Imported by: 0

Documentation

Index

Constants

View Source
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.

View Source
const (
	ErrParseError     = -32700
	ErrInvalidRequest = -32600
	ErrMethodNotFound = -32601
	ErrInvalidParams  = -32602
	ErrInternal       = -32603
)

JSON-RPC 2.0 error codes.

View Source
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

func IsConnectTimeout(err error) bool

IsConnectTimeout reports whether err was caused by a bounded IPC connect timeout.

func PeerPID

func PeerPID(ctx context.Context) int

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 Dial

func Dial(socketPath string) (*Client, error)

Dial connects to the IPC server at the given endpoint path.

func (*Client) Call

func (c *Client) Call(method string, params interface{}, result interface{}) error

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.

func (*Client) CallWithTimeout

func (c *Client) CallWithTimeout(method string, params interface{}, result interface{}, timeout time.Duration) error

CallWithTimeout is Call with a caller-selected read deadline.

func (*Client) Close

func (c *Client) Close() error

Close disconnects from the server.

type ConnectTimeoutError

type ConnectTimeoutError struct {
	SocketPath      string
	TimeoutDuration time.Duration
	Err             error
}

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 RPCError

type RPCError struct {
	Code    int    `json:"code"`
	Message string `json:"message"`
}

RPCError represents a JSON-RPC 2.0 error object.

func (*RPCError) Error

func (e *RPCError) Error() string

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

func NewRequest(method string, params interface{}) (*Request, error)

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

func NewErrorResponse(id int64, code int, message string) *Response

NewErrorResponse creates an error JSON-RPC 2.0 response.

func NewResponse

func NewResponse(id int64, result interface{}) (*Response, error)

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 NewServer

func NewServer() *Server

NewServer creates a new IPC server.

func (*Server) Close

func (s *Server) Close()

Close gracefully shuts down the server.

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

func (s *Server) Listen(socketPath string) error

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

func (s *Server) Serve(socketPath string) error

Serve starts listening on the given IPC endpoint path. It blocks until Close is called, then returns nil.

func (*Server) ServeReady

func (s *Server) ServeReady() error

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

type StreamFunc func(send func(interface{}) error) error

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.

Jump to

Keyboard shortcuts

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