Documentation
¶
Overview ¶
Package busctl glues the consume client to the bus daemon. It encapsulates the three operations a consumer needs at startup and shutdown:
discover: find the running bus or fork a fresh one (race-free, plan §12
P3 "try dial → try fork lock → retry dial")
spawn: exec `dws event _bus --client-id <id>` as a detached
background process (stdio detach, setsid, ready pipe handshake)
stop: gracefully terminate the bus daemon (SIGTERM + IPC fallback)
All three operations are short-lived helpers — they own no long-lived goroutines and return whole errors to their caller.
Index ¶
- Constants
- Variables
- func Discover(cfg DiscoverConfig) (net.Conn, error)
- func LockPath(workDir string) string
- func MetaPath(workDir string) string
- func QueryStatus(endpoint string) (*transport.StatusResp, error)
- func ReadyFDFromEnv() *os.File
- func Spawn(cfg SpawnConfig) (pid int, err error)
- func Stop(cfg StopConfig) error
- func StopConsumers(endpoint string, subscribeIDs []string) (transport.ConsumerStopResp, error)
- type BusEntry
- type BusEntryState
- type DiscoverConfig
- type EntryStatus
- type SpawnConfig
- type SpawnFunc
- type StopConfig
Constants ¶
const DefaultStatusRPCTimeout = 2 * time.Second
DefaultStatusRPCTimeout caps how long QueryStatus waits for the bus to reply. 2s is generous — the bus's status_resp is a synchronous in-memory snapshot, sub-millisecond in practice; the timeout exists only to bound pathological cases (bus stuck in shutdown).
const DefaultStopTimeout = 5 * time.Second
DefaultStopTimeout is the wall-clock budget Stop waits for each graceful or fallback termination phase. 5s covers the bus's own tear-down (broadcast Bye → consumer goroutines drain → cleanup) with margin.
const ReadyFDEnv = "DWS_EVENT_BUS_READY_FD"
ReadyFDEnv is the env var the spawned `event _bus` child inspects to find the ready-pipe write end. The parent passes the FD number; child opens it via os.NewFile(fd, "ready") and writes 'R' on success or 'E' on failure. 3 is the first FD slot beyond stdio in cmd.ExtraFiles.
Variables ¶
var ErrConsumerStopUnsupported = errors.New("busctl: targeted consumer stop is unsupported")
ErrConsumerStopUnsupported lets callers fall back to the legacy process signal when an already-running bus predates the targeted stop protocol.
var ErrNotRunning = errors.New("busctl: bus is not running")
ErrNotRunning indicates bus.lock either does not exist or its recorded PID is not alive. Stop returns this as a sentinel so the caller can distinguish "nothing to stop" from "failed to stop".
var ErrOwnerUnverified = errors.New("busctl: bus process ownership could not be verified")
ErrOwnerUnverified indicates the PID from bus.lock is alive but no longer owns that lock. Stop refuses to send a process-level signal in this state because the operating system may have reused the stale PID.
var ErrSpawnFailed = errors.New("busctl: bus child reported startup failure")
ErrSpawnFailed is returned when the child reports startup failure through the Unix ready pipe or exits before binding the Windows named pipe.
var ErrSpawnTimeout = errors.New("busctl: bus child did not signal readiness within deadline")
ErrSpawnTimeout is returned when ReadyTimeout elapses without any signal.
var ReadyTimeout = 10 * time.Second
ReadyTimeout caps how long Spawn waits for the child to signal readiness. 10s is generous — bus startup is local-only work (file I/O + socket bind), so 1s would normally suffice; the extra headroom covers cold-start keychain prompts and slow CI machines.
Functions ¶
func Discover ¶
func Discover(cfg DiscoverConfig) (net.Conn, error)
Discover returns a connected net.Conn to the bus for cfg.ClientID. If the bus is not running, Discover forks a new one and waits for it to come up.
Race-free three-step (plan §12 P3):
- try dial IPC endpoint → success → done
- failed → call Spawn (fork _bus); Spawn blocks until ready pipe says 'R' (or returns ErrSpawnFailed if another process won the race and our bus startup hit ErrBusy on the lock)
- dial again — should succeed; retry with backoff up to DialDeadline in case Spawn succeeded but socket bind has tiny latency
On concurrent Discover by N processes: only one Spawn succeeds (the others get ErrBusy via the bus daemon's own lock acquisition). Losers fall through to the retry-dial loop in step 3 and connect to the bus the winner brought up.
func LockPath ¶
LockPath returns the canonical bus.lock path for the given working dir. Centralised so consume / status / stop all agree on the location.
func QueryStatus ¶
func QueryStatus(endpoint string) (*transport.StatusResp, error)
QueryStatus dials the bus IPC, sends Hello with Role=status, sends a StatusReq, reads exactly one StatusResp, and closes. Returns the decoded response or an error if any step fails.
Used by `dws event status` and `dws event list` to fetch live per-consumer / per-event-type counters. The bus's handleStatusRPC path (see internal/event/bus/daemon.go) handles this connection without registering with the Hub — ad-hoc tooling does not count as a consumer in `status.active_consumers`.
func ReadyFDFromEnv ¶
ReadyFDFromEnv returns the inherited ready pipe (or nil if not set). The `event _bus` command handler calls this at startup, passes the returned *os.File to bus.Run as Config.ReadyPipe, and the bus signals readiness through it.
func Spawn ¶
func Spawn(cfg SpawnConfig) (pid int, err error)
Spawn starts the detached bus and waits for its inherited ready pipe. ExtraFiles is intentionally confined to Unix: Go does not support it on Windows.
func Stop ¶
func Stop(cfg StopConfig) error
Stop signals the bus daemon for cfg.WorkDir to exit gracefully and waits for the process to actually die. Returns ErrNotRunning if no bus is running for that work dir.
Stop first asks the bus to shut down through its owner-only IPC endpoint. This is the normal path on every platform and lets the daemon cancel its cloud source, notify consumers, and release its lock. If the endpoint is unavailable (for example an older bus) or graceful shutdown times out, the platform signal is used as a compatibility fallback: SIGTERM on Unix and TerminateProcess via os.Kill on Windows.
func StopConsumers ¶ added in v1.0.55
func StopConsumers(endpoint string, subscribeIDs []string) (transport.ConsumerStopResp, error)
StopConsumers asks a running bus to close only consumers matching the supplied personal subscription IDs.
Types ¶
type BusEntry ¶
type BusEntry struct {
WorkDir string `json:"workdir"`
Edition string `json:"edition"`
SourceKind dwsevent.SourceKind `json:"source_kind,omitempty"`
ClientIDHash string `json:"client_id_hash"`
IdentityHash string `json:"identity_hash,omitempty"`
HolderPID int `json:"holder_pid"`
State BusEntryState `json:"state"`
// Meta, if non-nil, lets list/status display the original ClientID
// (reverse-mapped from the hash) and the bus start time.
Meta *bus.Meta `json:"meta,omitempty"`
}
BusEntry is one bus working directory found on disk plus its detected lifecycle state. EnumerateBuses produces these; the cobra layer joins them with QueryStatus output to render the full status view.
func EnumerateBuses ¶
EnumerateBuses scans <configDir>/events/<editionFilter>/*/ for bus working directories. An empty editionFilter scans every edition directory found under events/.
Returns a deterministic slice sorted by (edition, source_kind, identity_hash). Missing/inaccessible directories are skipped silently — list/status commands should still succeed when only some editions have ever run a bus.
func FindBusByClientID ¶
FindBusByClientID is the "current ClientID" lookup used by `event status` (no --all). Returns the entry for the given (edition, clientIDHash) pair or nil if no bus has ever started for it. The caller derives the hash using event.ClientIDHash.
func FindBusByIdentity ¶
func FindBusByIdentity(configDir, editionName string, sourceKind dwsevent.SourceKind, identityHash string) *BusEntry
FindBusByIdentity looks up a bus in the source-kind-aware layout.
func (BusEntry) IPCEndpoint ¶
IPCEndpoint returns the IPC endpoint for this entry. Delegates to dwsevent.IPCEndpoint so status/stop dial exactly where consume and the bus daemon bound (a private per-user runtime path on Unix and a named pipe on Windows).
type BusEntryState ¶
type BusEntryState string
BusEntryState classifies a discovered bus directory's runtime state. Used by `dws event status/list` to render the table and by --fail-on-orphan to drive exit code.
const ( // BusStateRunning: bus.lock holds an alive PID — the daemon is up. BusStateRunning BusEntryState = "running" // BusStateOrphan: bus.meta exists but bus.lock PID is dead. The user // should `dws event stop --client-id <id>` (which detects the dead // PID and unblocks fresh starts) or rm -rf the working directory. BusStateOrphan BusEntryState = "orphan" // BusStateNotRunning: directory exists (e.g. bus.meta retained for // historic reasons) but bus.lock is missing or empty. Clean state. BusStateNotRunning BusEntryState = "not_running" )
type DiscoverConfig ¶
type DiscoverConfig struct {
WorkDir string
IPCEndpoint string
ClientID string
// Spawn is the fork-bus callback. Default busctl.Spawn.
Spawn SpawnFunc
// SpawnExtraArgs is forwarded to Spawn (for tests).
SpawnExtraArgs []string
// DialBackoff: initial sleep between retry dials when another process
// is spawning. Doubled each attempt up to DialMaxBackoff.
DialBackoff time.Duration
// DialMaxBackoff caps backoff.
DialMaxBackoff time.Duration
// DialDeadline caps total wall-clock time spent discovering.
DialDeadline time.Duration
}
DiscoverConfig describes one discover attempt. WorkDir holds bus.lock and persistent bus metadata; Unix sockets live in a private per-user runtime directory so WorkDir may reside on a shared filesystem without socket support. The caller must mkdir WorkDir with pkg/config.DirPerm beforehand.
type EntryStatus ¶
type EntryStatus struct {
Entry BusEntry `json:"entry"`
Live *transport.StatusResp `json:"live,omitempty"`
}
EntryStatus combines static FS info (BusEntry) with the live RPC snapshot (StatusResp). For not_running / orphan entries Live is nil.
func QueryEntry ¶
func QueryEntry(entry BusEntry) EntryStatus
QueryEntry fetches the live status for one BusEntry. Returns the entry wrapped with a nil Live when state != running (or when the dial fails). Errors from QueryStatus are folded into Live=nil so the caller's table rendering does not need to surface per-bus dial failures (they are already conveyed by State).
type SpawnConfig ¶
type SpawnConfig struct {
// ExecPath is the dws binary to exec. Default os.Executable().
ExecPath string
// ClientID is passed as `--client-id` to `dws event _bus`. Required.
ClientID string
// IPCEndpoint is the bus endpoint the child binds. Windows uses the
// existing named pipe as its readiness handshake because os/exec does not
// support ExtraFiles there. Unix keeps the inherited ready-pipe protocol.
IPCEndpoint string
// ExtraArgs are appended after `--client-id`. Empty for normal use; tests
// pass `--extra-flag-for-test` etc.
ExtraArgs []string
// Env to pass to the child. Defaults to os.Environ(). Unix appends the
// ReadyFDEnv entry automatically; Windows removes any inherited copy.
Env []string
}
SpawnConfig describes one spawn attempt.
type SpawnFunc ¶
type SpawnFunc func(SpawnConfig) (pid int, err error)
SpawnFunc abstracts the fork-bus operation so tests can inject a fake instead of execing a real binary. Production callers pass busctl.Spawn.
type StopConfig ¶
type StopConfig struct {
// WorkDir holds bus.lock; Stop reads the PID from there.
WorkDir string
// IPCEndpoint enables the cross-platform graceful stop RPC. Callers that
// know the bus identity should always provide it. Empty preserves the
// legacy signal-only fallback for compatibility and focused tests.
IPCEndpoint string
// Timeout is the wall-clock budget for graceful exit and, if required,
// the subsequent platform termination fallback.
Timeout time.Duration
}
StopConfig identifies the target bus and tunes timing.