Files
bee/audit/internal/app/blackbox.go
T
mchusandClaude Sonnet 5 8a91f0f783 fix(webui): repair broken scenario Run button onclick, dedupe build.sh overlay staging
- page_scenario.go: onclick built via JSON.stringify() embedded raw double
  quotes inside a double-quoted HTML attribute, truncating the attribute so
  the click handler never compiled; pass the name through an escaped
  data-scenario-name attribute instead.
- build.sh: overlay staging rsyncs (OVERLAY_DIR->stage, stage->includes.chroot)
  ran without --delete, so a scenario removed from the repo (a9924b0) stayed
  baked into every ISO built from the persistent stage cache since — the
  "second script" in the Scenario page's list.
- blackbox: rewritten around a deterministic local zip + incremental
  patch-the-changed-suffix onto removable media, instead of walking/copying
  ~90 files through a synchronous ntfs-3g FUSE mount every cycle. journalctl
  captures are now "--since last sync" (were "--since boot", growing with
  uptime) and metrics.db is excluded (was copied whole every cycle).
- scenario: nvbandwidth-acs-ab now escalates GPU count (same-socket pair,
  other socket's pair, one cross-socket pair, all GPUs) under each ACS state
  instead of always running all 6 GPUs at once, using a new `bee
  gpu-bandwidth-groups` subcommand that discovers socket layout from
  `nvidia-smi topo -m` at runtime — gpu_indices is host-specific, so this
  can't be baked into the scenario file.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-07-29 18:11:24 +03:00

966 lines
29 KiB
Go

package app
import (
"bytes"
"context"
"crypto/rand"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"os"
"os/exec"
"path/filepath"
"sort"
"strings"
"sync"
"syscall"
"time"
"bee/audit/internal/platform"
)
const (
blackboxMarkerName = ".bee-blackbox"
blackboxDiscoverInterval = 2 * time.Second
blackboxMinFlushPeriod = 1 * time.Second
blackboxMaxFlushPeriod = 30 * time.Second
blackboxRecoveryFastCount = 5
// blackboxKickFileName is a marker under exportDir that SAT job
// execution touches after each job finishes (see platform.SetJobBoundaryHook
// wiring in New()). Workers poll its mtime so a just-finished job's
// output reaches removable media promptly instead of waiting out
// whatever the current adaptive flush period happens to be.
blackboxKickFileName = ".blackbox-kick"
// blackboxSyncBracketTimeout bounds how long a syncBracket-marked SAT job
// (see platform.SetSyncBracketHook) blocks waiting for blackbox to catch
// up before/after the actual load step. Generous relative to the normal
// sub-second kick-to-sync latency, but bounded so a genuinely stuck or
// unplugged target can't hang the diagnostic indefinitely.
blackboxSyncBracketTimeout = 20 * time.Second
)
// blackboxKickPollInterval is how often an idle worker checks the kick file
// for a pending out-of-band sync request. A package var (not a const) so
// tests can shrink it instead of waiting out the real interval.
var blackboxKickPollInterval = 250 * time.Millisecond
// requestBlackboxSync touches the kick file under exportDir, signalling any
// running blackbox worker to sync on its next poll instead of waiting out
// its current flush period. Best-effort: a failure here just means the next
// scheduled sync picks up the data instead of an early one.
func requestBlackboxSync(exportDir string) {
path := filepath.Join(exportDir, blackboxKickFileName)
now := blackboxNow()
if err := os.Chtimes(path, now, now); err == nil {
return
}
_ = os.WriteFile(path, []byte(now.Format(time.RFC3339Nano)+"\n"), 0644)
}
// blackboxSyncWaitPollInterval is how often requestBlackboxSyncAndWait
// re-reads the state file while waiting for a sync to catch up. A package
// var so tests can shrink it.
var blackboxSyncWaitPollInterval = 200 * time.Millisecond
// requestBlackboxSyncAndWait kicks every enrolled blackbox target and blocks
// until each one has completed a sync cycle that started at or after this
// call, or until timeout elapses. Use this (instead of the fire-and-forget
// requestBlackboxSync) around an event where losing the data to a crash
// before it reaches removable media would matter — e.g. bracket a risky load
// step so its own SAT job directory (and the fact that it started at all) is
// durable on the flash drive before the load runs, and durable again once it
// finishes.
//
// Returns nil once caught up. Returns an error naming what didn't catch up
// (timeout, or a target stuck "degraded") — callers should treat that as
// best-effort informational (log it) rather than fail the job over it: a
// blackbox problem should not block the diagnostic the operator actually
// asked for.
func requestBlackboxSyncAndWait(exportDir, statePath string, timeout time.Duration) error {
requestedAt := blackboxNow()
requestBlackboxSync(exportDir)
deadline := time.Now().Add(timeout)
for {
state, err := ReadBlackboxState(statePath)
if err == nil {
if ok, pending := blackboxStateCaughtUp(state, requestedAt); ok {
return nil
} else if time.Now().After(deadline) {
return fmt.Errorf("blackbox sync did not catch up within %s (pending: %s)", timeout, pending)
}
} else if time.Now().After(deadline) {
return fmt.Errorf("blackbox sync wait: could not read state after %s: %w", timeout, err)
}
time.Sleep(blackboxSyncWaitPollInterval)
}
}
// blackboxStateCaughtUp reports whether every non-degraded enrolled target
// has synced at or after requestedAt. No targets enrolled counts as caught
// up (nothing to wait for — e.g. no removable media plugged in). A target
// stuck "degraded" (mount/copy failing) is skipped rather than waited on
// forever; its enrollment ID is named in the returned pending string so
// callers can log which target is the reason a real wait timed out.
func blackboxStateCaughtUp(state BlackboxState, requestedAt time.Time) (bool, string) {
for _, t := range state.Targets {
if t.Status == "degraded" {
continue
}
syncedAt, err := time.Parse(time.RFC3339Nano, t.LastSyncAtUTC)
if err != nil || syncedAt.Before(requestedAt) {
return false, t.EnrollmentID
}
}
return true, ""
}
// blackboxKickModTime returns the kick file's mtime, or the zero Time if it
// doesn't exist yet (nothing has requested a sync since boot).
func blackboxKickModTime(exportDir string) time.Time {
info, err := os.Stat(filepath.Join(exportDir, blackboxKickFileName))
if err != nil {
return time.Time{}
}
return info.ModTime()
}
var DefaultBlackboxStatePath = DefaultExportDir + "/blackbox-state.json"
var (
blackboxExecCommand = exec.Command
blackboxNow = func() time.Time { return time.Now().UTC() }
)
type BlackboxMarker struct {
Version int `json:"version"`
EnrollmentID string `json:"enrollment_id"`
CreatedAtUTC string `json:"created_at_utc"`
Host string `json:"host,omitempty"`
}
type BlackboxTargetStatus struct {
EnrollmentID string `json:"enrollment_id"`
Device string `json:"device"`
FS platform.RemovableTarget `json:"fs"`
BootFolder string `json:"boot_folder"`
Status string `json:"status"`
LastSyncAtUTC string `json:"last_sync_at_utc,omitempty"`
LastCycleDuration string `json:"last_cycle_duration,omitempty"`
FlushPeriod string `json:"flush_period"`
LastError string `json:"last_error,omitempty"`
Mountpoint string `json:"mountpoint,omitempty"`
}
type BlackboxState struct {
Status string `json:"status"`
BootStartedAtUTC string `json:"boot_started_at_utc"`
BootFolder string `json:"boot_folder"`
UpdatedAtUTC string `json:"updated_at_utc"`
Targets []BlackboxTargetStatus `json:"targets"`
}
type blackboxRuntime struct {
exportDir string
statePath string
system *platform.System
bootStarted time.Time
bootFolder string
mu sync.Mutex
workers map[string]*blackboxWorker
}
type discoveredBlackboxTarget struct {
marker BlackboxMarker
target platform.RemovableTarget
seenMount string
mountedByBee bool
}
type blackboxWorker struct {
runtime *blackboxRuntime
enrollmentID string
mu sync.Mutex
target platform.RemovableTarget
marker BlackboxMarker
mountpoint string
mountedByBee bool
status string
lastSyncAt time.Time
lastDuration time.Duration
flushPeriod time.Duration
lastError string
fastCycles int
lastKickSeen time.Time
stopCh chan struct{}
stoppedCh chan struct{}
// syncCycleFunc defaults to w.syncCycle; overridable in tests so run()'s
// wait/wake-on-kick logic can be exercised without a real removable-media
// mount and copy.
syncCycleFunc func() error
}
func RunBlackbox(ctx context.Context, exportDir, statePath string, system *platform.System) error {
exportDir = strings.TrimSpace(exportDir)
if exportDir == "" {
exportDir = DefaultExportDir
}
statePath = strings.TrimSpace(statePath)
if statePath == "" {
statePath = DefaultBlackboxStatePath
}
if system == nil {
system = platform.New()
}
bootStarted, err := bootStartedAtUTC()
if err != nil {
bootStarted = blackboxNow()
}
rt := &blackboxRuntime{
exportDir: exportDir,
statePath: statePath,
system: system,
bootStarted: bootStarted,
bootFolder: SupportBundleBaseName(bootStarted),
workers: make(map[string]*blackboxWorker),
}
_ = os.MkdirAll(filepath.Dir(statePath), 0755)
rt.persistState()
ticker := time.NewTicker(blackboxDiscoverInterval)
defer ticker.Stop()
for {
rt.reconcile()
select {
case <-ctx.Done():
rt.stopAll()
return ctx.Err()
case <-ticker.C:
}
}
}
func ReadBlackboxState(path string) (BlackboxState, error) {
path = strings.TrimSpace(path)
if path == "" {
path = DefaultBlackboxStatePath
}
raw, err := os.ReadFile(path)
if err != nil {
return BlackboxState{}, err
}
var state BlackboxState
if err := json.Unmarshal(raw, &state); err != nil {
return BlackboxState{}, err
}
return state, nil
}
func EnableBlackboxTarget(target platform.RemovableTarget) (BlackboxMarker, error) {
target = sanitizeRemovableTarget(target)
if target.Device == "" {
return BlackboxMarker{}, fmt.Errorf("device is required")
}
mountpoint, mountedByBee, err := ensureMountedTarget(target, "marker")
if err != nil {
return BlackboxMarker{}, err
}
defer func() {
if mountedByBee {
_ = unmountTarget(mountpoint)
}
}()
marker, _, err := readBlackboxMarker(mountpoint)
if err != nil && !errors.Is(err, os.ErrNotExist) {
return BlackboxMarker{}, err
}
if marker.EnrollmentID == "" {
marker = BlackboxMarker{
Version: 1,
EnrollmentID: newBlackboxEnrollmentID(),
CreatedAtUTC: blackboxNow().Format(time.RFC3339),
Host: hostnameOr("unknown"),
}
}
if err := writeBlackboxMarker(mountpoint, marker); err != nil {
return BlackboxMarker{}, err
}
return marker, nil
}
func DisableBlackboxTarget(device, enrollmentID string) error {
device = strings.TrimSpace(device)
enrollmentID = strings.TrimSpace(enrollmentID)
if device == "" && enrollmentID == "" {
return fmt.Errorf("device or enrollment_id is required")
}
system := platform.New()
targets, err := system.ListRemovableTargets()
if err != nil {
return err
}
for _, target := range targets {
target = sanitizeRemovableTarget(target)
mountpoint, mountedByBee, mountErr := ensureMountedTarget(target, "marker")
if mountErr != nil {
continue
}
remove := false
marker, _, err := readBlackboxMarker(mountpoint)
if err == nil {
if enrollmentID != "" && marker.EnrollmentID == enrollmentID {
remove = true
}
if device != "" && target.Device == device {
remove = true
}
}
if remove {
err = os.Remove(filepath.Join(mountpoint, blackboxMarkerName))
}
if mountedByBee {
_ = unmountTarget(mountpoint)
}
if remove {
return err
}
}
return os.ErrNotExist
}
func (rt *blackboxRuntime) reconcile() {
discovered, _ := rt.discoverMarkedTargets(rt.trackedDevices())
rt.mu.Lock()
defer rt.mu.Unlock()
seen := make(map[string]struct{}, len(discovered))
for _, found := range discovered {
seen[found.marker.EnrollmentID] = struct{}{}
worker, ok := rt.workers[found.marker.EnrollmentID]
if !ok {
worker = newBlackboxWorker(rt, found)
rt.workers[found.marker.EnrollmentID] = worker
go worker.run()
continue
}
worker.update(found)
}
for id, worker := range rt.workers {
if _, ok := seen[id]; ok {
continue
}
worker.stop()
delete(rt.workers, id)
}
rt.persistStateLocked()
}
func (rt *blackboxRuntime) stopAll() {
rt.mu.Lock()
workers := make([]*blackboxWorker, 0, len(rt.workers))
for _, worker := range rt.workers {
workers = append(workers, worker)
}
rt.workers = map[string]*blackboxWorker{}
rt.persistStateLocked()
rt.mu.Unlock()
for _, worker := range workers {
worker.stop()
}
}
// trackedDevices returns the already-enrolled targets, keyed by device path,
// for every worker currently running. Used by discoverMarkedTargets to skip
// re-mounting devices a worker already owns.
func (rt *blackboxRuntime) trackedDevices() map[string]discoveredBlackboxTarget {
rt.mu.Lock()
defer rt.mu.Unlock()
known := make(map[string]discoveredBlackboxTarget, len(rt.workers))
for _, worker := range rt.workers {
worker.mu.Lock()
if worker.target.Device != "" {
known[worker.target.Device] = discoveredBlackboxTarget{
marker: worker.marker,
target: worker.target,
}
}
worker.mu.Unlock()
}
return known
}
// discoverMarkedTargets probes removable media for the bee-blackbox marker.
// known holds devices a worker is already running against (see
// trackedDevices) — those are reported back as-is, without mounting, since
// mounting/unmounting the same device on every discovery tick (every
// blackboxDiscoverInterval) fights the worker's own mount for the device and
// was observed to slow its actual sync cycle by an order of magnitude on
// FUSE-backed filesystems (NTFS via ntfs-3g). Only devices with no running
// worker get probed.
func (rt *blackboxRuntime) discoverMarkedTargets(known map[string]discoveredBlackboxTarget) ([]discoveredBlackboxTarget, error) {
targets, err := rt.system.ListRemovableTargets()
if err != nil {
return nil, err
}
var out []discoveredBlackboxTarget
for _, rawTarget := range targets {
target := sanitizeRemovableTarget(rawTarget)
if target.Device == "" {
continue
}
if cached, ok := known[target.Device]; ok {
cached.target = target
out = append(out, cached)
continue
}
mountpoint, mountedByBee, err := ensureMountedTarget(target, "probe")
if err != nil {
continue
}
marker, ok, err := readBlackboxMarker(mountpoint)
if mountedByBee && !ok {
_ = unmountTarget(mountpoint)
}
if err != nil || !ok || marker.EnrollmentID == "" {
continue
}
if mountedByBee {
_ = unmountTarget(mountpoint)
}
out = append(out, discoveredBlackboxTarget{
marker: marker,
target: target,
seenMount: mountpoint,
mountedByBee: mountedByBee,
})
}
sort.Slice(out, func(i, j int) bool {
return out[i].marker.EnrollmentID < out[j].marker.EnrollmentID
})
return out, nil
}
func newBlackboxWorker(rt *blackboxRuntime, found discoveredBlackboxTarget) *blackboxWorker {
w := &blackboxWorker{
runtime: rt,
enrollmentID: found.marker.EnrollmentID,
target: found.target,
marker: found.marker,
flushPeriod: blackboxMinFlushPeriod,
status: "running",
stopCh: make(chan struct{}),
stoppedCh: make(chan struct{}),
}
w.syncCycleFunc = w.syncCycle
return w
}
func (w *blackboxWorker) run() {
defer close(w.stoppedCh)
kickPoll := time.NewTicker(blackboxKickPollInterval)
defer kickPoll.Stop()
for {
start := time.Now()
err := w.syncCycleFunc()
duration := time.Since(start)
w.finishCycle(duration, err)
w.lastKickSeen = blackboxKickModTime(w.runtime.exportDir)
wait := w.currentFlushPeriod()
timer := time.NewTimer(wait)
waitLoop:
for {
select {
case <-w.stopCh:
timer.Stop()
w.cleanup()
return
case <-timer.C:
break waitLoop
case <-kickPoll.C:
if mtime := blackboxKickModTime(w.runtime.exportDir); mtime.After(w.lastKickSeen) {
timer.Stop()
break waitLoop
}
}
}
}
}
func (w *blackboxWorker) update(found discoveredBlackboxTarget) {
w.mu.Lock()
defer w.mu.Unlock()
w.target = found.target
w.marker = found.marker
}
func (w *blackboxWorker) stop() {
select {
case <-w.stopCh:
default:
close(w.stopCh)
}
<-w.stoppedCh
}
func (w *blackboxWorker) currentFlushPeriod() time.Duration {
w.mu.Lock()
defer w.mu.Unlock()
return w.flushPeriod
}
func (w *blackboxWorker) finishCycle(duration time.Duration, err error) {
w.mu.Lock()
w.lastDuration = duration
if err != nil {
w.status = "degraded"
w.lastError = err.Error()
w.fastCycles = 0
w.flushPeriod = adjustFlushPeriod(w.flushPeriod, duration, false, 0)
} else {
w.status = "running"
w.lastSyncAt = blackboxNow()
w.lastError = ""
if duration <= w.flushPeriod/2 {
w.fastCycles++
} else {
w.fastCycles = 0
}
w.flushPeriod = adjustFlushPeriod(w.flushPeriod, duration, true, w.fastCycles)
}
w.mu.Unlock()
// persistState must be called without w.mu held: it acquires rt.mu then
// each worker.mu inside persistStateLocked, so holding w.mu here would
// cause a deadlock (w.mu → rt.mu → w.mu).
w.runtime.persistState()
}
func adjustFlushPeriod(current, duration time.Duration, success bool, fastCycles int) time.Duration {
if current <= 0 {
current = blackboxMinFlushPeriod
}
if duration <= 0 {
duration = current
}
next := current
if duration > current {
growA := time.Duration(float64(current) * 1.25)
growB := time.Duration(float64(duration) * 1.25)
if growB > growA {
next = growB
} else {
next = growA
}
}
if success && fastCycles >= blackboxRecoveryFastCount {
next = time.Duration(float64(current) * 0.9)
}
if next < blackboxMinFlushPeriod {
next = blackboxMinFlushPeriod
}
if next > blackboxMaxFlushPeriod {
next = blackboxMaxFlushPeriod
}
return next
}
// syncCycle stages the current export tree on fast local storage (never
// touching the removable-media mountpoint for the expensive part), packs it
// into a single deterministic zip, and patches only the changed suffix of
// that zip onto the target device — see blackbox_archive.go for why: walking
// and rewriting ~90 small files through a synchronous FUSE mount (ntfs-3g
// -o sync) was taking ~2x the flush period, dominated by journalctl output
// that grows with uptime and got fully re-read/re-written every cycle.
func (w *blackboxWorker) syncCycle() error {
target, marker := w.snapshotTarget()
mountpoint, mountedByBee, err := ensureMountedTarget(target, marker.EnrollmentID)
if err != nil {
return err
}
w.recordMountpoint(mountpoint, mountedByBee)
stageRoot := filepath.Join(w.runtime.exportDir, ".blackbox-stage", w.enrollmentID)
if err := os.RemoveAll(stageRoot); err != nil {
return err
}
// includeMetricsDB=false: metrics.db is a live, growing SQLite file: past
// versions copied it whole every cycle, which alone could dominate a
// cycle's cost as it grew. Not carried onto the blackbox mirror.
if err := categorizeExportTree(w.runtime.exportDir, stageRoot, false); err != nil {
return err
}
if err := w.captureSnapshots(stageRoot, w.lastCaptureSince()); err != nil {
return err
}
// Same doc pair the support bundle ships at its root — a blackbox
// capture on removable media has no manifest.txt/support-bundle
// equivalent to explain its layout, so without this an agent handed
// only the media would have nothing pointing it at README.md.
if err := writeBundleDocs(stageRoot); err != nil {
return err
}
cacheDir := filepath.Join(w.runtime.exportDir, ".blackbox-cache")
if err := os.MkdirAll(cacheDir, 0755); err != nil {
return err
}
cachedPath := filepath.Join(cacheDir, w.enrollmentID+".zip")
newZipPath := filepath.Join(cacheDir, w.enrollmentID+".zip.new")
if err := buildZipArchive(stageRoot, newZipPath); err != nil {
return err
}
targetPath := filepath.Join(mountpoint, w.runtime.bootFolder+" blackbox.zip")
if err := patchArchiveOnTarget(targetPath, newZipPath, cachedPath); err != nil {
return err
}
if err := os.Rename(newZipPath, cachedPath); err != nil {
return err
}
return nil
}
// lastCaptureSince is the "--since" boundary for this cycle's incremental
// journalctl captures: the previous successful sync, or boot time on the
// very first cycle. Using the last sync instead of always "--since boot"
// keeps each cycle's journalctl output bounded by the flush period instead
// of growing with total uptime.
func (w *blackboxWorker) lastCaptureSince() time.Time {
w.mu.Lock()
defer w.mu.Unlock()
if w.lastSyncAt.IsZero() {
return w.runtime.bootStarted
}
return w.lastSyncAt
}
func (w *blackboxWorker) cleanup() {
w.mu.Lock()
mountpoint := w.mountpoint
mountedByBee := w.mountedByBee
w.mu.Unlock()
if mountedByBee && mountpoint != "" {
_ = unmountTarget(mountpoint)
}
}
func (w *blackboxWorker) snapshotTarget() (platform.RemovableTarget, BlackboxMarker) {
w.mu.Lock()
defer w.mu.Unlock()
return w.target, w.marker
}
func (w *blackboxWorker) recordMountpoint(mountpoint string, mountedByBee bool) {
w.mu.Lock()
defer w.mu.Unlock()
w.mountpoint = mountpoint
w.mountedByBee = mountedByBee
}
// captureSnapshots writes this cycle's journalctl/dmesg/status snapshots
// into root (a local staging tree — see syncCycle). journalctl output is
// captured incrementally ("--since" the last successful cycle, not boot) and
// appended as a new timestamped entry under journal-increments/ each cycle,
// rather than overwriting one ever-growing file: re-dumping "--since boot"
// every cycle made each cycle's journalctl call (and the file it produced)
// grow with total uptime, independent of how much actually happened since
// the last sync.
func (w *blackboxWorker) captureSnapshots(root string, since time.Time) error {
cycleTS := blackboxNow().Format("20060102-150405.000")
sinceArg := since.Format(time.RFC3339)
incDir := filepath.Join(root, "tasks", "_services", "journal-increments")
if err := captureCommandAtomic(filepath.Join(incDir, "combined-"+cycleTS+".log"), "journalctl", "--no-pager", "--since", sinceArg); err != nil {
return err
}
for _, svc := range supportBundleServices {
dir := filepath.Join(root, serviceBundleDir(svc))
if err := captureCommandAtomic(filepath.Join(incDir, svc+"-"+cycleTS+".log"), "journalctl", "--no-pager", "-u", svc, "--since", sinceArg); err != nil {
return err
}
if err := captureCommandAtomic(filepath.Join(dir, svc+".status.txt"), "systemctl", "status", svc, "--no-pager"); err != nil {
return err
}
}
if err := captureCommandAtomic(filepath.Join(root, "livecd", "host", "dmesg.txt"), "dmesg"); err != nil {
return err
}
for _, item := range supportBundleOptionalFiles {
if err := copyFileIfChanged(item.src, filepath.Join(root, item.name)); err != nil && !errors.Is(err, os.ErrNotExist) {
return err
}
}
return nil
}
func (rt *blackboxRuntime) persistState() {
rt.mu.Lock()
defer rt.mu.Unlock()
rt.persistStateLocked()
}
func (rt *blackboxRuntime) persistStateLocked() {
state := BlackboxState{
Status: "disabled",
BootStartedAtUTC: rt.bootStarted.Format(time.RFC3339),
BootFolder: rt.bootFolder,
UpdatedAtUTC: blackboxNow().Format(time.RFC3339),
Targets: make([]BlackboxTargetStatus, 0, len(rt.workers)),
}
if len(rt.workers) > 0 {
state.Status = "running"
}
for _, worker := range rt.workers {
worker.mu.Lock()
targetState := BlackboxTargetStatus{
EnrollmentID: worker.enrollmentID,
Device: worker.target.Device,
FS: worker.target,
BootFolder: rt.bootFolder,
Status: worker.status,
FlushPeriod: worker.flushPeriod.String(),
LastError: worker.lastError,
Mountpoint: worker.mountpoint,
}
if !worker.lastSyncAt.IsZero() {
// Nanosecond precision, not RFC3339's default seconds — a sync
// that lands in the same wall-clock second as a kick request
// must still compare as "after" it (see blackboxStateCaughtUp),
// which second-only precision can get wrong.
targetState.LastSyncAtUTC = worker.lastSyncAt.Format(time.RFC3339Nano)
}
if worker.lastDuration > 0 {
targetState.LastCycleDuration = worker.lastDuration.String()
}
if worker.status == "degraded" {
state.Status = "degraded"
}
worker.mu.Unlock()
state.Targets = append(state.Targets, targetState)
}
sort.Slice(state.Targets, func(i, j int) bool {
return state.Targets[i].EnrollmentID < state.Targets[j].EnrollmentID
})
_ = writeJSONAtomic(rt.statePath, state)
}
func bootStartedAtUTC() (time.Time, error) {
raw, err := os.ReadFile("/proc/stat")
if err != nil {
return time.Time{}, err
}
for _, line := range strings.Split(string(raw), "\n") {
line = strings.TrimSpace(line)
if !strings.HasPrefix(line, "btime ") {
continue
}
parts := strings.Fields(line)
if len(parts) != 2 {
break
}
sec, err := time.ParseDuration(parts[1] + "s")
if err != nil {
break
}
return time.Unix(int64(sec/time.Second), 0).UTC(), nil
}
return time.Time{}, fmt.Errorf("boot time not found")
}
func newBlackboxEnrollmentID() string {
var buf [8]byte
if _, err := rand.Read(buf[:]); err != nil {
return fmt.Sprintf("bb-%d", time.Now().UnixNano())
}
return "bb-" + hex.EncodeToString(buf[:])
}
func sanitizeRemovableTarget(target platform.RemovableTarget) platform.RemovableTarget {
target.Device = strings.TrimSpace(target.Device)
target.FSType = strings.TrimSpace(target.FSType)
target.Size = strings.TrimSpace(target.Size)
target.Label = strings.TrimSpace(target.Label)
target.Model = strings.TrimSpace(target.Model)
target.Mountpoint = strings.TrimSpace(target.Mountpoint)
return target
}
func ensureMountedTarget(target platform.RemovableTarget, suffix string) (mountpoint string, mountedByBee bool, retErr error) {
target = sanitizeRemovableTarget(target)
if target.Mountpoint != "" {
if err := ensureWritableBlackboxMountpoint(target.Mountpoint); err == nil {
return target.Mountpoint, false, nil
}
}
mountpoint = filepath.Join("/tmp", "bee-blackbox-"+sanitizeFilename(suffix))
if err := os.MkdirAll(mountpoint, 0755); err != nil {
return "", false, err
}
// -o sync makes every write to the target synchronous at the VFS layer
// (no page-cache write-back to lose on a hard reset) instead of relying
// solely on the explicit syscall.Sync() calls in writeFileAtomic/
// unmountTarget to flush it after the fact. Those explicit syncs stay in
// place as a fallback (e.g. for target.Mountpoint above, an
// already-mounted filesystem we don't control the options of) — with -o
// sync already doing the work, they become a fast no-op most of the time
// instead of the primary durability mechanism.
if raw, err := blackboxExecCommand("mount", "-o", "sync", target.Device, mountpoint).CombinedOutput(); err != nil {
return "", false, formatBlackboxMountTargetError(target, string(raw), err)
}
if err := ensureWritableBlackboxMountpoint(mountpoint); err != nil {
_ = unmountTarget(mountpoint)
return "", false, err
}
return mountpoint, true, nil
}
func unmountTarget(mountpoint string) error {
syscall.Sync()
raw, err := blackboxExecCommand("umount", mountpoint).CombinedOutput()
if err != nil {
msg := strings.TrimSpace(string(raw))
if msg == "" {
return err
}
return fmt.Errorf("%s: %w", msg, err)
}
return nil
}
func readBlackboxMarker(mountpoint string) (BlackboxMarker, bool, error) {
raw, err := os.ReadFile(filepath.Join(mountpoint, blackboxMarkerName))
if err != nil {
if errors.Is(err, os.ErrNotExist) {
return BlackboxMarker{}, false, os.ErrNotExist
}
return BlackboxMarker{}, false, err
}
var marker BlackboxMarker
if err := json.Unmarshal(raw, &marker); err != nil {
return BlackboxMarker{}, false, err
}
return marker, true, nil
}
func writeBlackboxMarker(mountpoint string, marker BlackboxMarker) error {
if marker.Version == 0 {
marker.Version = 1
}
return writeJSONAtomic(filepath.Join(mountpoint, blackboxMarkerName), marker)
}
func copyFileIfChanged(src, dst string) error {
info, err := os.Stat(src)
if err != nil {
return err
}
if info.IsDir() {
return os.MkdirAll(dst, info.Mode().Perm())
}
srcData, err := os.ReadFile(src)
if err != nil {
return err
}
if dstData, err := os.ReadFile(dst); err == nil && bytes.Equal(dstData, srcData) {
return nil
}
return writeFileAtomic(dst, srcData, info.Mode().Perm())
}
func captureCommandAtomic(dst string, name string, args ...string) error {
raw, err := blackboxExecCommand(name, args...).CombinedOutput()
if len(raw) == 0 {
if err != nil {
raw = []byte(err.Error() + "\n")
} else {
raw = []byte("no output\n")
}
}
return writeFileAtomic(dst, raw, 0644)
}
func writeJSONAtomic(path string, v any) error {
raw, err := json.MarshalIndent(v, "", " ")
if err != nil {
return err
}
raw = append(raw, '\n')
return writeFileAtomic(path, raw, 0644)
}
func writeFileAtomic(path string, data []byte, perm os.FileMode) error {
if err := os.MkdirAll(filepath.Dir(path), 0755); err != nil {
return err
}
if existing, err := os.ReadFile(path); err == nil && bytes.Equal(existing, data) {
return nil
}
tmp := path + ".tmp"
f, err := os.OpenFile(tmp, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, perm)
if err != nil {
return err
}
if _, err := f.Write(data); err != nil {
_ = f.Close()
return err
}
if err := f.Sync(); err != nil {
_ = f.Close()
return err
}
if err := f.Close(); err != nil {
return err
}
if err := os.Rename(tmp, path); err != nil {
return err
}
return syncFilesystem(filepath.Dir(path))
}
// syncFilesystem flushes pending writes to durable storage. Uses the sync(2)
// syscall directly rather than spawning /bin/sync — writeFileAtomic calls
// this once per copied file (syncDirectoryTree can touch hundreds of files
// per cycle), and syscall.Sync() does the same flush without a fork+exec per
// call. path is unused (sync(2) always flushes system-wide; kept as a
// parameter for call-site clarity about which tree the caller cares about).
func syncFilesystem(path string) error {
syscall.Sync()
return nil
}
func ensureWritableBlackboxMountpoint(mountpoint string) error {
probe, err := os.CreateTemp(mountpoint, ".bee-blackbox-write-test-*")
if err != nil {
return fmt.Errorf("target filesystem is not writable: %w", err)
}
name := probe.Name()
if closeErr := probe.Close(); closeErr != nil {
_ = os.Remove(name)
return closeErr
}
if err := os.Remove(name); err != nil {
return err
}
return nil
}
func formatBlackboxMountTargetError(target platform.RemovableTarget, raw string, err error) error {
msg := strings.TrimSpace(raw)
fstype := strings.ToLower(strings.TrimSpace(target.FSType))
if fstype == "exfat" && strings.Contains(strings.ToLower(msg), "unknown filesystem type 'exfat'") {
return fmt.Errorf("mount %s: exFAT support is missing in this ISO build: %w", target.Device, err)
}
if msg == "" {
return err
}
return fmt.Errorf("%s: %w", msg, err)
}