v0.144.1 code: the wrapper survives a dead reader (BrokenPipe) so a killed pass still keeps its report; the agent looks for kept copies every 5 min (R-868, measured live)
gates / gates (push) Successful in 20s
gates / gates (push) Successful in 20s
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_0159rPz1ZhFKsS53msqPYxtS
This commit is contained in:
@@ -870,11 +870,10 @@ func runDaemon(cfg config.Config, logger *slog.Logger, logRing *applog.Ring) int
|
|||||||
"interval_s", hcfg.PollSeconds) // hub key intentionally not logged
|
"interval_s", hcfg.PollSeconds) // hub key intentionally not logged
|
||||||
// R-868 (v0.144.0): an OS pass whose agent was killed kept its report on disk — send it now (a pass that starts
|
// R-868 (v0.144.0): an OS pass whose agent was killed kept its report on disk — send it now (a pass that starts
|
||||||
// first sends it itself; the pass lock keeps the two apart).
|
// first sends it itself; the pass lock keeps the two apart).
|
||||||
go func() {
|
// v0.144.1: and again every 5 minutes — the wrapper of a killed pass can finish AFTER the restart (measured).
|
||||||
if n := osLeg.SendUnsent(ctx); n > 0 {
|
go osLeg.SendUnsentLoop(ctx, 5*time.Minute, func(n int) {
|
||||||
logger.Info("osupdate: sent kept report(s) at start", "count", n)
|
logger.Info("osupdate: sent kept report(s)", "count", n)
|
||||||
}
|
})
|
||||||
}()
|
|
||||||
|
|
||||||
// Reconcile (slice 4) runs alongside the hub loop, sharing the per-guest queue
|
// Reconcile (slice 4) runs alongside the hub loop, sharing the per-guest queue
|
||||||
// (doc 03 §10). At slice 4 the desired-state provider is empty (no hub serving
|
// (doc 03 §10). At slice 4 the desired-state provider is empty (no hub serving
|
||||||
|
|||||||
+12
-1
@@ -263,7 +263,14 @@ class Runner:
|
|||||||
os.replace(tmp, NONCE_FILE)
|
os.replace(tmp, NONCE_FILE)
|
||||||
|
|
||||||
def log(self, line):
|
def log(self, line):
|
||||||
|
# R-868 (v0.144.1): the agent that reads stderr may be GONE (killed mid-pass, measured live on demo-hp
|
||||||
|
# 2026-10-05): the write then raises BrokenPipeError, and v0.144.0 died right there — after apt had installed
|
||||||
|
# everything, before the report copy was saved. A dead reader must never stop the pass; the journal still
|
||||||
|
# gets every line. Pinned by AgentDiesMidPass.
|
||||||
|
try:
|
||||||
print(line, file=sys.stderr, flush=True)
|
print(line, file=sys.stderr, flush=True)
|
||||||
|
except OSError:
|
||||||
|
pass
|
||||||
try:
|
try:
|
||||||
subprocess.run(["logger", "-t", "felhom-os-apply", line], timeout=10)
|
subprocess.run(["logger", "-t", "felhom-os-apply", line], timeout=10)
|
||||||
except Exception:
|
except Exception:
|
||||||
@@ -1453,7 +1460,11 @@ def main(argv, runner=None, environ=None):
|
|||||||
r.save_report(argv[2], a.report)
|
r.save_report(argv[2], a.report)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
r.log(f"os-apply: the report copy could not be saved (the run is unaffected): {e}")
|
r.log(f"os-apply: the report copy could not be saved (the run is unaffected): {e}")
|
||||||
print("OSAPPLY-REPORT " + json.dumps(a.report, sort_keys=True))
|
try:
|
||||||
|
print("OSAPPLY-REPORT " + json.dumps(a.report, sort_keys=True), flush=True)
|
||||||
|
except OSError:
|
||||||
|
# R-868: nobody reads stdout any more (the agent was killed); the copy above carries the report.
|
||||||
|
sys.stdout = open(os.devnull, "w")
|
||||||
return rc
|
return rc
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -1077,3 +1077,49 @@ class SaveReportOnDisk(unittest.TestCase):
|
|||||||
with self.assertRaises(OSError):
|
with self.assertRaises(OSError):
|
||||||
self.r.save_report(os.path.join(link, "plan-r3-guest-apply.json"), {"mode": "apply"})
|
self.r.save_report(os.path.join(link, "plan-r3-guest-apply.json"), {"mode": "apply"})
|
||||||
self.assertEqual(os.listdir(real), [])
|
self.assertEqual(os.listdir(real), [])
|
||||||
|
|
||||||
|
|
||||||
|
class AgentDiesMidPass(unittest.TestCase):
|
||||||
|
"""R-868, MEASURED LIVE 2026-10-05 05:45 UTC on demo-hp (agent v0.144.0): the agent was kill -9-ed while apt-get
|
||||||
|
ran; apt finished all 13 packages, but the wrapper's next log line went to a stderr pipe nobody reads any more ->
|
||||||
|
BrokenPipeError -> the wrapper died before save_report: no DONE in the journal, no kept copy, no report.
|
||||||
|
COMPANION RED-PROOF: let Runner.log write to stderr unguarded -> this test fails (BrokenPipeError)."""
|
||||||
|
|
||||||
|
def test_a_dead_reader_does_not_stop_the_report_copy(self):
|
||||||
|
import io, sys, contextlib
|
||||||
|
|
||||||
|
class DeadPipe(io.TextIOBase):
|
||||||
|
dead = False
|
||||||
|
def write(self, s):
|
||||||
|
if DeadPipe.dead:
|
||||||
|
raise BrokenPipeError(32, "Broken pipe")
|
||||||
|
return len(s)
|
||||||
|
def flush(self):
|
||||||
|
if DeadPipe.dead:
|
||||||
|
raise BrokenPipeError(32, "Broken pipe")
|
||||||
|
|
||||||
|
class DyingFake(Fake):
|
||||||
|
def emulate(self, argv):
|
||||||
|
a = [x for x in argv if not re.match(r"^[A-Z_]+=", x) and x != "env"]
|
||||||
|
if a and a[0] == "apt-get" and "install" in a and "-s" not in a and "--print-uris" not in a:
|
||||||
|
DeadPipe.dead = True # the agent is killed while apt-get runs
|
||||||
|
return super().emulate(argv)
|
||||||
|
|
||||||
|
f = DyingFake()
|
||||||
|
journal = []
|
||||||
|
f.log = lambda line: osapply.Runner.log(f, line) # the REAL log(): stderr, then the journal
|
||||||
|
prev_run = osapply.subprocess.run
|
||||||
|
osapply.subprocess.run = lambda argv, *a, **kw: (journal.append(argv[-1]) if argv[0] == "logger"
|
||||||
|
else prev_run(argv, *a, **kw)) # never the real `logger`
|
||||||
|
prev_err, prev_out = sys.stderr, sys.stdout
|
||||||
|
sys.stderr, sys.stdout = DeadPipe(), DeadPipe()
|
||||||
|
try:
|
||||||
|
rc = osapply.main(["felhom-os-apply", "--plan", PLAN], runner=f)
|
||||||
|
finally:
|
||||||
|
sys.stderr, sys.stdout = prev_err, prev_out
|
||||||
|
osapply.subprocess.run = prev_run
|
||||||
|
DeadPipe.dead = False
|
||||||
|
self.assertEqual(rc, 0)
|
||||||
|
self.assertEqual(len(f.saved_reports), 1, "the kept copy is the only way this report reaches the hub")
|
||||||
|
self.assertEqual(len(f.saved_reports[0][1]["upgraded"]), 2)
|
||||||
|
self.assertTrue(any(l.startswith("os-apply: DONE") for l in journal), "the journal must still get the DONE line")
|
||||||
|
|||||||
@@ -7,6 +7,7 @@ import (
|
|||||||
"path/filepath"
|
"path/filepath"
|
||||||
"strings"
|
"strings"
|
||||||
"syscall"
|
"syscall"
|
||||||
|
"time"
|
||||||
|
|
||||||
"gitea.dooplex.hu/admin/felhom-agent/internal/hub"
|
"gitea.dooplex.hu/admin/felhom-agent/internal/hub"
|
||||||
)
|
)
|
||||||
@@ -62,6 +63,24 @@ func (l *Leg) SendUnsent(ctx context.Context) int {
|
|||||||
return l.sendUnsentLocked(ctx)
|
return l.sendUnsentLocked(ctx)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// SendUnsentLoop runs SendUnsent now and then every `every` until ctx ends (v0.144.1). MEASURED live on demo-hp
|
||||||
|
// 2026-10-05: after a kill -9 the daemon restarted in ~5 s, while the orphaned root wrapper was still installing — its
|
||||||
|
// copy appeared ~7 s AFTER the start-time sender had looked. One look at start is therefore not enough. A glob of the
|
||||||
|
// plan dir every few minutes costs nothing; the pass lock keeps it off a running pass. Pinned by
|
||||||
|
// TestR868_ACopyWrittenAfterTheStartIsSentByTheLoop.
|
||||||
|
func (l *Leg) SendUnsentLoop(ctx context.Context, every time.Duration, onSent func(int)) {
|
||||||
|
for {
|
||||||
|
if n := l.SendUnsent(ctx); n > 0 && onSent != nil {
|
||||||
|
onSent(n)
|
||||||
|
}
|
||||||
|
select {
|
||||||
|
case <-ctx.Done():
|
||||||
|
return
|
||||||
|
case <-time.After(every):
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func (l *Leg) sendUnsentLocked(ctx context.Context) int {
|
func (l *Leg) sendUnsentLocked(ctx context.Context) int {
|
||||||
if l.Hub == nil {
|
if l.Hub == nil {
|
||||||
return 0 // nobody to send to: keep the copies for a process that has the hub
|
return 0 // nobody to send to: keep the copies for a process that has the hub
|
||||||
|
|||||||
@@ -7,6 +7,7 @@ import (
|
|||||||
"os"
|
"os"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
"testing"
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
"gitea.dooplex.hu/admin/felhom-agent/internal/hub"
|
"gitea.dooplex.hu/admin/felhom-agent/internal/hub"
|
||||||
)
|
)
|
||||||
@@ -115,6 +116,29 @@ func TestR868_RunningPassKeepsTheSenderOff(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// v0.144.1 — THE LIVE SHAPE (demo-hp 2026-10-05 05:45 UTC): the restarted daemon looked at 05:45:09, the orphaned
|
||||||
|
// wrapper wrote its copy at ~05:45:16. The loop must still send it.
|
||||||
|
// COMPANION RED-PROOF: make SendUnsentLoop return after the first look → "the late copy was never sent".
|
||||||
|
func TestR868_ACopyWrittenAfterTheStartIsSentByTheLoop(t *testing.T) {
|
||||||
|
w := &fakeWrapper{t: t}
|
||||||
|
l, h := newLeg(t, w, nil)
|
||||||
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
|
defer cancel()
|
||||||
|
sent := make(chan int, 4)
|
||||||
|
go l.SendUnsentLoop(ctx, 20*time.Millisecond, func(n int) { sent <- n })
|
||||||
|
time.Sleep(50 * time.Millisecond) // the start-time look found nothing
|
||||||
|
must(t, os.WriteFile(reportFile(l.PlanDir, "late", LayerGuest, "apply"),
|
||||||
|
[]byte(`{"mode":"apply","layer":"guest","run_id":"late","trigger":"debug","ring":0,"vmid":9201,"upgraded":[{"name":"openssl","version":"u3"}]}`), 0o600))
|
||||||
|
select {
|
||||||
|
case n := <-sent:
|
||||||
|
if n != 1 || len(h.reports) != 1 || h.reports[0].RunID != "late" || h.reports[0].Outcome == "" {
|
||||||
|
t.Fatalf("sent %d: %+v", n, h.reports)
|
||||||
|
}
|
||||||
|
case <-time.After(3 * time.Second):
|
||||||
|
t.Fatal("the late copy was never sent")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func must(t *testing.T, err error) {
|
func must(t *testing.T, err error) {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
Reference in New Issue
Block a user