Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
31 changes: 31 additions & 0 deletions docs/operations/metrics.md
Original file line number Diff line number Diff line change
Expand Up @@ -132,6 +132,8 @@ counter, not evidence that the kernel still enforces every packet.
| `executor_schedule_remaining_seconds` | Minimum nonnegative remaining signing lifetime, using the announced start, epoch length and chain length on the dispatcher's clock. |
| `executor_clock_estimated_error_seconds` | Maximum reported kernel error estimate; this is not a measured uncertainty bound. |
| `executor_disclosure_held_seconds` | Maximum reported age of an installed-key disclosure hold. It does not measure delivery or durable storage of keys at the dispatcher. |
| `executors_disclosure_completion_unknown` | Executors without a fresh sender schedule-clock completion sample. |
| `executor_disclosure_completion_upper_bound_seconds` | Maximum executor-reported current-generation time from the oldest newly acknowledged key becoming due to receipt of its durable-storage acknowledgment, including the reply path. |
| `executors_disclosure_lag_unknown` | Executors whose disclosure delivery lag cannot be determined. |
| `executor_disclosure_lag_seconds` | Maximum disclosure delivery lag: time since the oldest due key became disclosable without this dispatcher having verified and recorded it; zero when every due key of each executor's current chain is stored. |

Expand Down Expand Up @@ -246,3 +248,32 @@ failures have their own fixed reasons. Host filesystem calls have no imposed
kernel deadline; the configured database directory must be on a supported
local filesystem. Non-Linux host measurements are explicitly unsupported.
Available disk bytes are an observation, not a promise a future write succeeds.

### Sender disclosure completion

`executor_disclosure_completion_upper_bound_seconds` complements the dispatcher
backlog gauge. After verifying and committing a key, the dispatcher acknowledges
coverage through a requested epoch. A duplicate or later committed key can cover
earlier epochs without another database write. The acknowledgment establishes
that verified durable completion happened before the reply; it does not promise
indefinite archive retention. Failed writes, rejected disclosures and memory-only
stores produce no durable receipt. Replies are bounded to five requested chains.

The executor measures from the scheduled due time of the oldest newly covered
positive epoch to receipt of the acknowledgment, using that current generation's
local schedule clock. This includes holder/scheduling delay, delivery, verification,
storage and the return path. Lost replies and retries can increase the upper bound;
skipped epochs are included. The high-water mark survives control reconnects, but
duplicate replies do not renew a sample. The separate observation is unavailable
before a receipt, after a failed call or missing receipt, and once its actual sample
age reaches one minute. Missing samples are not zero. Fresh report receipt alone
does not refresh the sample's age.

This is an executor-reported local-clock completion bound, not exact commit time,
one-way latency, a receiver retrieval measurement or independently calibrated clock
uncertainty. Linux monotonic clocks can exclude suspend. Observed clock-unready,
degraded clock health or wall/monotonic drift invalidates the observation for that
generation, even if the clock later recovers. Recovered chains have no original
monotonic reading and remain timing-unknown; their final keys are retained until a
matching durable receipt arrives or the existing retention limit expires. Older
dispatchers return no receipts, so this observation remains unknown with them.
110 changes: 110 additions & 0 deletions internal/dispatcher/disclosure_receipt_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,110 @@
// SPDX-License-Identifier: Apache-2.0
// Copyright 2026 ETH Zurich

package dispatcher

import (
"bytes"
"testing"
"time"

"github.com/netsec-ethz/debuglet/internal/dispatcher/tag"
pb "github.com/netsec-ethz/debuglet/protocol"
)

type heldDisclosureBackend struct {
tag.Backend
entered, release chan struct{}
}

func (b *heldDisclosureBackend) Save(id string, chain tag.Chain, epoch int64, key []byte, at time.Time) error {
close(b.entered)
<-b.release
return b.Backend.Save(id, chain, epoch, key, at)
}

func TestDisclosureReceiptFollowsRealDurableCommit(t *testing.T) {
d, db, _ := newRegistryFixture(t)
const id = "receipt"
tail := bytes.Repeat([]byte{0x6D}, 32)
anchor := chainKey(tail, 0)
attributionRegister(t, d, id, anchor)
held := &heldDisclosureBackend{Backend: attributionBackend{db: db}, entered: make(chan struct{}), release: make(chan struct{})}
d.keystore = tag.NewPersistentKeyStore(held)
mutation := effectTestMutation(t, d, id)
type result struct {
response *pb.HeartbeatResponse
err error
}
done := make(chan result, 1)
go func() {
defer mutation.Finish()
response, err := d.OnHeartbeat(mutation.Context(), mutation, &pb.HeartbeatRequest{ExecutorId: id, TeslaKeyEpoch: 3, TeslaKey: chainKey(tail, 3)})
done <- result{response, err}
}()
select {
case <-held.entered:
case <-time.After(5 * time.Second):
t.Fatal("write not reached")
}
select {
case got := <-done:
t.Fatal("receipt before commit", got)
default:
}
if n := attributionKeyRows(t, d); n != 0 {
t.Fatal("uncommitted row visible", n)
}
close(held.release)
got := <-done
if got.err != nil || len(got.response.GetDisclosureReceipts()) != 1 || got.response.DisclosureReceipts[0].StoredThroughEpoch != 3 {
t.Fatal(got)
}
if n := attributionKeyRows(t, d); n != 1 {
t.Fatal("receipt before row commit", n)
}
// A restarted key store proves duplicate coverage from the actual database.
d.keystore = tag.NewPersistentKeyStore(attributionBackend{db: db})
heartbeat := func(epoch int64, key []byte, extra []*pb.TeslaDisclosure) []*pb.TeslaDisclosureReceipt {
t.Helper()
m := effectTestMutation(t, d, id)
defer m.Finish()
out, err := d.OnHeartbeat(m.Context(), m, &pb.HeartbeatRequest{ExecutorId: id, TeslaKeyEpoch: epoch, TeslaKey: key, ExtraDisclosures: extra})
if err != nil {
t.Fatal(err)
}
return out.DisclosureReceipts
}
if receipts := heartbeat(3, chainKey(tail, 3), nil); len(receipts) != 1 {
t.Fatal(receipts)
}
if _, err := db.Exec("CREATE TRIGGER fail_disclosure BEFORE INSERT ON attribution_keys BEGIN SELECT RAISE(FAIL, 'controlled storage failure'); END"); err != nil {
t.Fatal(err)
}
if receipts := heartbeat(5, tail, nil); len(receipts) != 0 {
t.Fatal("failed storage acknowledged", receipts)
}
if n := attributionKeyRows(t, d); n != 1 {
t.Fatal(n)
}
if _, err := db.Exec("DROP TRIGGER fail_disclosure"); err != nil {
t.Fatal(err)
}
if receipts := heartbeat(5, tail, nil); len(receipts) != 1 || receipts[0].StoredThroughEpoch != 5 {
t.Fatal(receipts)
}
// An ignored lower submission is a prefix receipt, not acceptance of its bytes.
if receipts := heartbeat(4, bytes.Repeat([]byte{0xFF}, 32), nil); len(receipts) != 1 || receipts[0].StoredThroughEpoch != 4 {
t.Fatal(receipts)
}
if receipts := heartbeat(5, bytes.Repeat([]byte{0xFF}, 32), nil); len(receipts) != 0 {
t.Fatal("conflicting key acknowledged", receipts)
}
if receipts := heartbeat(0, nil, make([]*pb.TeslaDisclosure, 5)); len(receipts) != 0 {
t.Fatal("oversized extras acknowledged", receipts)
}
d.keystore = tag.NewKeyStore()
if receipts := heartbeat(5, tail, nil); len(receipts) != 0 {
t.Fatal("memory store acknowledged", receipts)
}
}
51 changes: 51 additions & 0 deletions internal/dispatcher/metrics_delivery_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,51 @@
// SPDX-License-Identifier: Apache-2.0
// Copyright 2026 ETH Zurich

package dispatcher

import (
"bytes"
"testing"
"time"

pb "github.com/netsec-ethz/debuglet/protocol"
)

func TestMetricsDisclosureCompletionExpiresActualSample(t *testing.T) {
now := time.Now()
anchor := bytes.Repeat([]byte{1}, 32)
report := &pb.VantagePointReport{SchemaVersion: 1, DisclosureDelivery: &pb.DisclosureDeliveryObservation{Anchor: anchor, StoredThroughEpoch: 3, ScheduledToAckNs: int64(2 * time.Second), SampleAgeNs: int64(30 * time.Second)}}
entry := &executorEntry{RegisteredExecutor: &RegisteredExecutor{TeslaAnchorKey: bytes.Clone(anchor), TeslaChainLength: 10}}
entry.Capabilities = capabilitiesFromReport(&pb.ExecutorCapabilities{SchemaVersion: 1, Attribution: &pb.AttributionState{State: "available"}}, now)
entry.capabilityObserved, entry.vantageObserved = now, now
entry.vantage = vantageFromReport(report)
observe := func(at time.Time) ExecutorResourceMetric {
var metrics ExecutorHealthMetrics
metrics.observe(entry, at, true)
return metrics.DisclosureCompletion
}
if got := observe(now.Add(29 * time.Second)); got.Unknown != 0 || got.Value == nil || *got.Value != 2 {
t.Fatal(got)
}
if got := observe(now.Add(30 * time.Second)); got.Unknown != 1 || got.Value != nil {
t.Fatal("sample survived exact boundary", got)
}
// Input is copied before the caller can mutate it.
report.DisclosureDelivery.Anchor[0] = 7
if got := observe(now); got.Unknown != 0 {
t.Fatal("report retained caller storage", got)
}
entry.TeslaAnchorKey = bytes.Repeat([]byte{9}, 32)
if got := observe(now); got.Unknown != 1 {
t.Fatal("other generation accepted", got)
}
entry.TeslaAnchorKey = anchor
report.DisclosureDelivery.SampleAgeNs = -1
if got := vantageFromReport(report); got.disclosureDelivery != nil {
t.Fatal("negative age accepted")
}
report.DisclosureDelivery.SampleAgeNs = int64(time.Minute)
if got := vantageFromReport(report); got.disclosureDelivery != nil {
t.Fatal("expired age accepted")
}
}
16 changes: 16 additions & 0 deletions internal/dispatcher/metrics_health.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,11 +4,13 @@
package dispatcher

import (
"bytes"
"math"
"time"

"github.com/netsec-ethz/debuglet/internal/dispatcher/tag"
"github.com/netsec-ethz/debuglet/internal/observability"
pb "github.com/netsec-ethz/debuglet/protocol"
)

// ExecutorHealthMetrics aggregates the registry's validated reports. Missing,
Expand Down Expand Up @@ -38,6 +40,7 @@ type ExecutorHealthMetrics struct {

// DisclosureLag is the maximum disclosure delivery lag in seconds.
DisclosureLag ExecutorResourceMetric
DisclosureCompletion ExecutorResourceMetric
RefusedAdmissions, RevokedSockets ExecutorResourceMetric
DropVerdicts, DropSKBBytes [2]ExecutorResourceMetric // Ingress, egress; current counter instance.
}
Expand Down Expand Up @@ -69,6 +72,19 @@ func (m *ExecutorHealthMetrics) observe(e *executorEntry, now time.Time, connect
resources = e.vantage.resources
}
m.observeResources(resources)
var completion *float64
if e.vantage != nil && fresh(e.vantageObserved) && e.vantage.disclosureDelivery != nil &&
e.Capabilities != nil && fresh(e.capabilityObserved) && e.Capabilities.Attribution != nil &&
e.Capabilities.Attribution.Reason != "clock_unready" && e.Capabilities.Attribution.Reason != "clock_drift" {
d := e.vantage.disclosureDelivery
age := now.Sub(e.vantageObserved)
if bytes.Equal(d.Anchor, e.TeslaAnchorKey) && d.StoredThroughEpoch < e.TeslaChainLength &&
time.Duration(d.SampleAgeNs) < pb.DisclosureSampleLifetime-age {
seconds := float64(d.ScheduledToAckNs) / float64(time.Second)
completion = &seconds
}
}
m.DisclosureCompletion.observe(completion, false)
var refused, revoked *float64
if e.vantage != nil && fresh(e.vantageObserved) && e.vantage.networkDenials != nil {
r, c := float64(e.vantage.networkDenials.RefusedAdmissions), float64(e.vantage.networkDenials.RevokedSockets)
Expand Down
27 changes: 20 additions & 7 deletions internal/dispatcher/rpc_handlers.go
Original file line number Diff line number Diff line change
Expand Up @@ -91,39 +91,44 @@ func (d *Dispatcher) OnHeartbeat(ctx context.Context, mutation *rpc.Mutation, re
})
entry.Write(zap.Int64("amount", earnings.TotalIncome), zap.Int64("next payout", earnings.CurrentBalance))
}
d.storeDisclosure(ctx, execID, chain, seen, req.GetTeslaKeyAnchor(), req.GetTeslaKeyEpoch(), req.GetTeslaKey())
response := &pb.HeartbeatResponse{}
if receipt := d.storeDisclosure(ctx, execID, chain, seen, req.GetTeslaKeyAnchor(), req.GetTeslaKeyEpoch(), req.GetTeslaKey()); receipt != nil {
response.DisclosureReceipts = append(response.DisclosureReceipts, receipt)
}
// Keys of earlier chains, such as the tail of the chain a restart retired,
// follow the same path. The keystore keeps at most maxExtraDisclosures
// chains per executor, so a longer list is refused whole.
if extra := req.GetExtraDisclosures(); len(extra) > maxExtraDisclosures {
d.logger.Warn("Refused a heartbeat's extra TESLA disclosures", zap.String("executor_id", daemonlog.Identifier(execID)), zap.Int("count", len(extra)))
} else {
for _, disclosure := range extra {
d.storeDisclosure(ctx, execID, chain, seen, disclosure.GetAnchor(), disclosure.GetEpoch(), disclosure.GetKey())
if receipt := d.storeDisclosure(ctx, execID, chain, seen, disclosure.GetAnchor(), disclosure.GetEpoch(), disclosure.GetKey()); receipt != nil {
response.DisclosureReceipts = append(response.DisclosureReceipts, receipt)
}
}
}
return &pb.HeartbeatResponse{}, nil
return response, nil
}

// maxExtraDisclosures bounds the extra disclosures one heartbeat may carry; it
// matches the chains the keystore keeps per executor.
const maxExtraDisclosures = 4
const maxExtraDisclosures = pb.MaxDisclosureReceipts - 1

// storeDisclosure verifies one disclosed key and stores it. An empty anchor
// names chain, the session's own. A key for another chain of this executor is
// verified against that chain's recorded schedule, so the tail of a chain can
// be disclosed after the executor restarted onto a new one. An anchor not on
// record names no chain the dispatcher knows, and its key is dropped.
func (d *Dispatcher) storeDisclosure(ctx context.Context, execID string, chain tag.Chain, seen time.Time, anchor []byte, epoch int64, key []byte) {
func (d *Dispatcher) storeDisclosure(ctx context.Context, execID string, chain tag.Chain, seen time.Time, anchor []byte, epoch int64, key []byte) *pb.TeslaDisclosureReceipt {
if len(anchor) > 0 && len(key) > 0 && !bytes.Equal(anchor, chain.Anchor) {
recorded, ok, err := d.recordedChain(ctx, execID, anchor)
if err != nil {
d.logger.Warn("Failed to read the recorded TESLA chain of a disclosure", zap.String("executor_id", execID), zap.Error(err))
return
return nil
}
if !ok {
d.logger.Debug("Dropped a disclosed TESLA key for a chain not on record", zap.String("executor_id", execID), zap.String("chain", tag.ChainID(anchor)))
return
return nil
}
chain = recorded
}
Expand All @@ -147,6 +152,14 @@ func (d *Dispatcher) storeDisclosure(ctx context.Context, execID string, chain t
// executor's next heartbeat discloses the key again.
d.logger.Warn("Failed to record disclosed TESLA key", zap.String("executor_id", execID), zap.Error(err))
}
if err != nil || epoch <= 0 || len(key) == 0 {
return nil
}
through, ok := d.keystore.DurableThrough(execID, chain.Anchor)
if !ok {
return nil
}
return &pb.TeslaDisclosureReceipt{Anchor: append([]byte(nil), chain.Anchor...), StoredThroughEpoch: min(epoch, through)}
}

func (d *Dispatcher) OnResources(ctx context.Context, mutation *rpc.Mutation, req *pb.ResourcesRequest) (*pb.ResourcesResponse, error) {
Expand Down
46 changes: 46 additions & 0 deletions internal/dispatcher/tag/disclosure_receipt_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,46 @@
// SPDX-License-Identifier: Apache-2.0
// Copyright 2026 ETH Zurich

package tag

import (
"errors"
"testing"
"time"
)

func TestDurableDisclosureCoverageRequiresCommittedBackend(t *testing.T) {
keys := hashChain(0xD3, 8)
chain := Chain{Anchor: keys[0]}
memory := NewKeyStore()
if err := memory.Store("executor", chain, time.Time{}, 3, keys[3]); err != nil {
t.Fatal(err)
}
if _, ok := memory.DurableThrough("executor", chain.Anchor); ok {
t.Fatal("memory-only store acknowledged durable completion")
}
backend := &memoryBackend{keys: map[string][]byte{}, fail: errors.New("write failure")}
durable := NewPersistentKeyStore(backend)
if err := durable.Store("executor", chain, time.Time{}, 3, keys[3]); err == nil {
t.Fatal("write succeeded")
}
if _, ok := durable.DurableThrough("executor", chain.Anchor); ok {
t.Fatal("failed write acknowledged")
}
backend.fail = nil
for _, epoch := range []int64{3, 3, 5, 4} {
if err := durable.Store("executor", chain, time.Time{}, epoch, keys[epoch]); err != nil {
t.Fatal(err)
}
}
if through, ok := durable.DurableThrough("executor", chain.Anchor); !ok || through != 5 || backend.saves != 2 {
t.Fatalf("through=%d ok=%v saves=%d", through, ok, backend.saves)
}
restarted := NewPersistentKeyStore(backend)
if err := restarted.Store("executor", chain, time.Time{}, 4, keys[4]); err != nil {
t.Fatal(err)
}
if through, ok := restarted.DurableThrough("executor", chain.Anchor); !ok || through != 5 {
t.Fatal(through, ok)
}
}
11 changes: 11 additions & 0 deletions internal/dispatcher/tag/keystore.go
Original file line number Diff line number Diff line change
Expand Up @@ -441,3 +441,14 @@ func (ks *KeyStore) PrintKeys() {
}
}
}

// DurableThrough returns cached verified coverage backed by a successful durable
// write (or read of that record). Memory-only stores cannot acknowledge durable
// delivery. It performs no I/O and a cache miss is unknown.
func (ks *KeyStore) DurableThrough(executorID string, anchor []byte) (int64, bool) {
if ks.backend == nil {
return 0, false
}
epoch, cached := ks.CachedLatest(executorID, anchor)
return epoch, cached && epoch > 0
}
2 changes: 2 additions & 0 deletions internal/dispatcher/transport/api/handlers_metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -147,6 +147,7 @@ func formatMetrics(control dispatcher.ControlMetrics, host observability.HostSna
}
gauge("executors_schedule_unknown", "Registered executors without a usable current schedule observation.", h.ScheduleUnknown)
gauge("executors_schedule_expired", "Registered executors whose announced signing schedule expired on the dispatcher clock.", h.ScheduleExpired)
gauge("executors_disclosure_completion_unknown", "Registered executors without a fresh sender schedule-clock durable-disclosure completion sample.", h.DisclosureCompletion.Unknown)
gauge("executors_disclosure_lag_unknown", "Registered executors whose disclosure delivery lag cannot be determined.", h.DisclosureLag.Unknown)
for _, metric := range []struct {
name, help string
Expand All @@ -156,6 +157,7 @@ func formatMetrics(control dispatcher.ControlMetrics, host observability.HostSna
{"executor_schedule_remaining_seconds", "Minimum remaining signing lifetime of announced executor schedules, using the dispatcher clock.", h.ScheduleRemainingSeconds, h.ScheduleUnknown},
{"executor_clock_estimated_error_seconds", "Maximum reported kernel clock error estimate; not an independently verified uncertainty bound.", h.ClockEstimatedErrorSeconds, h.ClockEstimateUnknown},
{"executor_disclosure_held_seconds", "Maximum reported installed-key disclosure hold age; not end-to-end key delivery lag.", &h.DisclosureHeldSeconds, h.AttributionUnknown},
{"executor_disclosure_completion_upper_bound_seconds", "Maximum executor-reported current-generation schedule-clock bound from disclosure due to durable acknowledgment, including the reply path; not calibrated one-way latency.", h.DisclosureCompletion.Value, h.DisclosureCompletion.Unknown},
{"executor_disclosure_lag_seconds", "Maximum time since the oldest due key became disclosable without this dispatcher verifying and recording it.", h.DisclosureLag.Value, h.DisclosureLag.Unknown},
} {
reason := ""
Expand Down
Loading