feat(cloud): stream windowed container metrics to Dozzle Cloud (#4911)

Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
Amir Raminfar
2026-08-15 10:16:30 -07:00
committed by GitHub
co-authored by Claude Opus 5
parent 8f45ef6ed6
commit 2b49998a60
26 changed files with 1881 additions and 183 deletions
+1
View File
@@ -147,6 +147,7 @@ declare module 'vue' {
'Mdi:poll': typeof import('~icons/mdi/poll')['default']
'Mdi:refresh': typeof import('~icons/mdi/refresh')['default']
'Mdi:satelliteVariant': typeof import('~icons/mdi/satellite-variant')['default']
'Mdi:shieldCheckOutline': typeof import('~icons/mdi/shield-check-outline')['default']
'Mdi:textBoxOutline': typeof import('~icons/mdi/text-box-outline')['default']
'Mdi:trashCanOutline': typeof import('~icons/mdi/trash-can-outline')['default']
'Mdi:webhook': typeof import('~icons/mdi/webhook')['default']
+50 -12
View File
@@ -83,19 +83,57 @@
</div>
</div>
<label class="flex min-h-13 cursor-pointer items-center justify-between gap-4 p-4">
<div class="flex flex-col gap-0.5">
<span class="text-sm font-medium">{{ $t("cloud.stream-logs") }}</span>
<span class="text-base-content/60 text-xs">{{ $t("cloud.stream-logs-help") }}</span>
<!--
One toggle gates BOTH log lines and container metrics — they ride the
same connection, and opting out of shipping log contents implies
opting out of shipping resource usage. The copy has to spell out
everything that leaves the instance, or the toggle understates itself.
-->
<div class="p-4">
<label class="flex min-h-13 cursor-pointer items-center justify-between gap-4">
<div class="flex flex-col gap-0.5">
<span class="text-sm font-medium">{{ $t("cloud.privacy.toggle") }}</span>
<span class="text-base-content/60 text-xs">{{ $t("cloud.privacy.required-for") }}</span>
</div>
<input
type="checkbox"
class="toggle toggle-primary toggle-sm shrink-0"
:checked="streamLogs"
:disabled="isSavingStreamLogs"
@change="onStreamLogsChange(($event.target as HTMLInputElement).checked)"
/>
</label>
<div class="border-base-content/10 mt-3 space-y-3 rounded-md border p-3 text-xs">
<div v-if="streamLogs">
<p class="text-base-content/70 font-medium">{{ $t("cloud.privacy.sends-heading") }}</p>
<ul class="text-base-content/60 mt-1.5 space-y-1">
<li class="flex items-start gap-1.5">
<mdi:text-box-outline class="mt-0.5 shrink-0 text-sm" />
<span>{{ $t("cloud.privacy.sends-logs") }}</span>
</li>
<li class="flex items-start gap-1.5">
<mdi:chart-line class="mt-0.5 shrink-0 text-sm" />
<span>{{ $t("cloud.privacy.sends-metrics") }}</span>
</li>
</ul>
</div>
<p v-else class="text-base-content/60">{{ $t("cloud.privacy.off-note") }}</p>
<div v-if="streamLogs" class="text-base-content/60 flex items-start gap-1.5">
<mdi:shield-check-outline class="text-success mt-0.5 shrink-0 text-sm" />
<span>
<span class="text-base-content/70 font-medium">{{ $t("cloud.privacy.never-heading") }}</span>
{{ $t("cloud.privacy.never-env") }}
</span>
</div>
<p v-if="streamLogs" class="text-base-content/45 border-base-content/10 border-t pt-2">
{{ $t("cloud.privacy.per-container") }}
</p>
</div>
<input
type="checkbox"
class="toggle toggle-primary toggle-sm shrink-0"
:checked="streamLogs"
:disabled="isSavingStreamLogs"
@change="onStreamLogsChange(($event.target as HTMLInputElement).checked)"
/>
</label>
</div>
<div class="flex gap-2 p-4">
<a :href="cloudUrl" target="_blank" rel="noreferrer noopener" class="btn btn-sm">
+21 -7
View File
@@ -24,13 +24,13 @@ import (
)
const (
initialBackoff = 1 * time.Second
maxBackoff = 30 * time.Second
backoffFactor = 2
jitterFraction = 0.1
maxConcurrent = 5
maxConcurrentStreams = 10
unauthenticatedPause = 1 * time.Hour
initialBackoff = 1 * time.Second
maxBackoff = 30 * time.Second
backoffFactor = 2
jitterFraction = 0.1
maxConcurrent = 5
maxConcurrentStreams = 10
unauthenticatedPause = 1 * time.Hour
)
// Client manages the gRPC connection to Dozzle Cloud
@@ -289,6 +289,20 @@ func (c *Client) connect(ctx context.Context, apiKey string) (wasConnected bool,
log.Debug().Msg("host service does not support log streaming; skipping")
}
// Container metrics ride the same connection and the same privacy toggle:
// a user who opted out of shipping log contents has opted out of shipping
// their containers' resource usage too.
if streamLogs {
if sshs, ok := c.deps.HostService.(StatsStreamHostService); ok {
stats := newStatsStreamer(sshs, c.deps.Labels, sendResp)
wg.Go(func() {
stats.run(streamLifetime)
})
} else {
log.Debug().Msg("host service does not support stats streaming; skipping")
}
}
defer func() {
// Cancel all active log streams before shutting down
c.activeStreams.Range(func(key, value any) bool {
+438
View File
@@ -0,0 +1,438 @@
package cloud
import (
"cmp"
"context"
"maps"
"slices"
"sync"
"time"
"github.com/amir20/dozzle/internal/container"
pb "github.com/amir20/dozzle/proto/cloud"
"github.com/rs/zerolog/log"
)
// StatSample pairs a raw container stat with the host it came from.
// container.ContainerStat carries only an ID, and the cloud host service fans
// stats in from several local client services, so the host has to be stamped
// at the fan-in point or it is lost.
type StatSample struct {
Stat container.ContainerStat
HostID string
}
// StatsStreamHostService is the subset of the host service the stats streamer
// needs. It is type-asserted at connect time; when the running host service
// doesn't satisfy it (k8s) the streamer is skipped, exactly as
// LogStreamHostService is today.
type StatsStreamHostService interface {
ToolHostService
SubscribeStats(ctx context.Context, samples chan<- StatSample)
}
const (
// statsWindow is the aggregation window. Docker pushes stats at roughly
// 1Hz, so this folds ~30 samples per container into one entry.
statsWindow = 30 * time.Second
// statsMaxEntries caps a single batch. A host with more running containers
// than this reports a truncated batch and logs the drop — never silently.
statsMaxEntries = 2000
// statsSampleChanBuf absorbs a tick's worth of samples from a busy host
// without blocking the collector's fan-out loop.
statsSampleChanBuf = 512
// statsMaxSendFailures is how many consecutive failed sends end the
// streamer. A wedged stream should not spin forever, but one bad send must
// not cost the connection every later window either.
statsMaxSendFailures = 5
// statsMaxStaleTicks is how many windows an accumulator may survive without
// its container resolving in the metadata snapshot before it is discarded.
// One tick of grace covers a container that started mid-window; beyond that
// it is gone and holding the samples only leaks memory.
statsMaxStaleTicks = 2
)
// seriesKey identifies an accumulator. Keyed on container ID because that is
// what a raw stat carries; the ID is resolved to a name at emit time and never
// leaves this process — see the StatsBatchEntry comment in cloud.proto for why
// the wire format is keyed on name instead.
type seriesKey struct {
hostID string
containerID string
}
// wireKey is how Cloud identifies a series: host plus container name, with no
// ID (see the StatsBatchEntry comment in cloud.proto). Two accumulators can
// collapse onto one of these when a container is replaced mid-window.
type wireKey struct {
hostID string
name string
}
// statAcc folds a window's samples for one container. Gauges accumulate a sum
// (for the mean) and a running max; counters are monotonic totals, so only the
// most recent value matters.
type statAcc struct {
cpuSum, cpuMax float64
memSum, memMax float64
memUsageSum float64
samples uint32
networkRx, networkTx uint64
diskRead, diskWrite uint64
staleTicks int
}
// containerMeta is the per-window metadata snapshot needed to turn raw stats
// into a wire entry: the container's name, the core count to normalise CPU by,
// and whether the user opted the container out.
type containerMeta struct {
name string
cores float64
disabled bool
}
// statsStreamer aggregates raw container stats into fixed windows and pushes
// one StatsBatch per window to Dozzle Cloud. Unlike logStreamer it runs a
// single goroutine for every container on the instance — stats already arrive
// multiplexed on one channel, so there is nothing to parallelise.
type statsStreamer struct {
hostService StatsStreamHostService
labels container.ContainerLabels
send func(resp *pb.ToolResponse) error
acc map[seriesKey]*statAcc
// meta is the last-good container metadata snapshot. It is owned entirely
// by run()'s goroutine: refreshes happen off-loop and hand the new snapshot
// back over metaCh, so there is no shared state and no lock.
meta map[seriesKey]containerMeta
// metaCh delivers a completed async refresh. Buffered so a refresh that
// finishes after run() has returned doesn't leak its goroutine.
metaCh chan map[seriesKey]containerMeta
// refreshing guards against piling up refreshes when the host service is
// slower than the window.
refreshing bool
// window is the flush interval. A field rather than a bare const so tests
// can drive real windows through run() instead of reaching into state the
// run goroutine owns.
window time.Duration
// listWarnOnce keeps a permanently unreachable host from warning every
// window. Touched only from buildMeta's goroutine chain.
listWarnOnce sync.Once
}
func newStatsStreamer(hostService StatsStreamHostService, labels container.ContainerLabels, send func(resp *pb.ToolResponse) error) *statsStreamer {
return &statsStreamer{
hostService: hostService,
labels: labels,
send: send,
acc: make(map[seriesKey]*statAcc),
meta: make(map[seriesKey]containerMeta),
metaCh: make(chan map[seriesKey]containerMeta, 1),
window: statsWindow,
}
}
// run blocks until ctx is cancelled, folding samples into per-container
// accumulators and flushing one batch per window.
func (ss *statsStreamer) run(ctx context.Context) {
samples := make(chan StatSample, statsSampleChanBuf)
ss.hostService.SubscribeStats(ctx, samples)
// Prime the metadata map so the first window can already resolve names
// rather than spending its grace tick on a cold cache. Async, like every
// later refresh — see startMetaRefresh for why this must never block.
ss.startMetaRefresh(ctx)
ticker := time.NewTicker(ss.window)
defer ticker.Stop()
log.Debug().Dur("window", ss.window).Msg("stats streamer: started")
defer log.Debug().Msg("stats streamer: stopped")
// Consecutive failed sends. Reset by any success — this is for a stream
// that is wedged rather than merely unlucky.
failures := 0
for {
select {
case <-ctx.Done():
return
case s := <-samples:
ss.observe(s)
case m := <-ss.metaCh:
ss.meta = m
ss.refreshing = false
case now := <-ticker.C:
// Kick the next refresh before flushing so the snapshot is as fresh
// as possible for the NEXT window; this one flushes against the
// last-good snapshot.
ss.startMetaRefresh(ctx)
if err := ss.flush(now); err != nil {
// A failed send is not proof the connection is gone, and this
// used to return — one transient error and the instance
// reported no metrics again until it reconnected, saying so
// only at Debug. A genuinely dead stream cancels ctx, which is
// the case above, so the honest response here is to log and
// keep windowing.
failures++
log.Warn().Err(err).Int("consecutive", failures).Msg("stats streamer: send failed")
if failures >= statsMaxSendFailures {
log.Error().Int("consecutive", failures).Msg("stats streamer: giving up until reconnect")
return
}
continue
}
failures = 0
}
}
}
// observe folds one raw sample into its accumulator. Samples for a container
// the user disabled are dropped here so they never occupy memory.
func (ss *statsStreamer) observe(s StatSample) {
key := seriesKey{hostID: s.HostID, containerID: s.Stat.ID}
if m, ok := ss.meta[key]; ok && m.disabled {
return
}
a, ok := ss.acc[key]
if !ok {
a = &statAcc{}
ss.acc[key] = a
}
a.cpuSum += s.Stat.CPUPercent
a.cpuMax = max(a.cpuMax, s.Stat.CPUPercent)
a.memSum += s.Stat.MemoryPercent
a.memMax = max(a.memMax, s.Stat.MemoryPercent)
a.memUsageSum += s.Stat.MemoryUsage
a.samples++
// Counters are cumulative since container start — keep the newest.
a.networkRx = s.Stat.NetworkRxTotal
a.networkTx = s.Stat.NetworkTxTotal
a.diskRead = s.Stat.DiskReadTotal
a.diskWrite = s.Stat.DiskWriteTotal
}
// flush converts every resolvable accumulator into a wire entry and pushes them
// as one batch, against the last-good metadata snapshot. Emitted and abandoned
// accumulators are removed, so a container that stops simply drops out of the
// next window.
//
// Deliberately does NO I/O: it runs on the same goroutine that drains the sample
// channel, and blocking here backs up into the stats collector (see
// startMetaRefresh).
func (ss *statsStreamer) flush(now time.Time) error {
// Keyed the way the wire is keyed, not the way the accumulators are: a
// redeploy inside one window leaves the dying and the starting container
// both holding samples under the same name, and emitting them as two
// entries would put two points on the same series at the same timestamp.
merged := make(map[wireKey]*pb.StatsBatchEntry, len(ss.acc))
tsNs := now.UnixNano()
for key, a := range ss.acc {
if a.samples == 0 {
delete(ss.acc, key)
continue
}
m, ok := ss.meta[key]
if !ok {
// Container hasn't shown up in a snapshot yet (started mid-window)
// or has already gone. Hold the samples for one more window, then
// give up rather than accumulate forever.
a.staleTicks++
if a.staleTicks >= statsMaxStaleTicks {
delete(ss.acc, key)
}
continue
}
delete(ss.acc, key)
if m.disabled {
continue
}
n := float64(a.samples)
cores := m.cores
if cores <= 0 {
cores = 1
}
entry := &pb.StatsBatchEntry{
HostId: key.hostID,
ContainerName: m.name,
TimestampNs: tsNs,
// Stat.CPUPercent is per-core (100% == one full core). Normalise by
// the core count so Cloud sees overall load, matching the UI and the
// metric-alert path in internal/notification/processing.go.
CpuPercent: a.cpuSum / n / cores,
CpuPercentMax: a.cpuMax / cores,
MemoryPercent: a.memSum / n,
MemoryPercentMax: a.memMax,
MemoryUsageBytes: a.memUsageSum / n,
NetworkRxTotal: a.networkRx,
NetworkTxTotal: a.networkTx,
DiskReadTotal: a.diskRead,
DiskWriteTotal: a.diskWrite,
Samples: a.samples,
CpuCores: cores,
}
wk := wireKey{hostID: key.hostID, name: m.name}
if prev, ok := merged[wk]; ok {
mergeEntry(prev, entry)
} else {
merged[wk] = entry
}
}
if len(merged) == 0 {
return nil
}
entries := slices.Collect(maps.Values(merged))
// Sort before any truncation so a host over the cap reports a stable subset
// rather than a different random slice every window.
slices.SortFunc(entries, func(a, b *pb.StatsBatchEntry) int {
if c := cmp.Compare(a.HostId, b.HostId); c != 0 {
return c
}
return cmp.Compare(a.ContainerName, b.ContainerName)
})
if len(entries) > statsMaxEntries {
log.Warn().
Int("total", len(entries)).
Int("sent", statsMaxEntries).
Msg("stats streamer: batch exceeds entry cap, truncating")
entries = entries[:statsMaxEntries]
}
return ss.send(&pb.ToolResponse{Type: &pb.ToolResponse_StatsBatch{StatsBatch: &pb.StatsBatch{Entries: entries}}})
}
// mergeEntry folds src into dst for two accumulators that resolved to the same
// (host, container name) — a container replaced mid-window. Means are
// re-weighted by sample count; maxima take the larger.
//
// Counters take the larger rather than the sum: Docker names are unique per
// host, so the two are never live at the same time, and the old container's
// cumulative total is the one the series was already sitting at. Summing would
// emit a total neither container ever reported. The next window drops back to
// the new container's counters, which reads as a counter reset — exactly what
// a restart is, and what rate() already knows how to handle.
func mergeEntry(dst, src *pb.StatsBatchEntry) {
total := float64(dst.Samples) + float64(src.Samples)
if total == 0 {
return
}
weighted := func(a, b float64) float64 {
return (a*float64(dst.Samples) + b*float64(src.Samples)) / total
}
dst.CpuPercent = weighted(dst.CpuPercent, src.CpuPercent)
dst.MemoryPercent = weighted(dst.MemoryPercent, src.MemoryPercent)
dst.MemoryUsageBytes = weighted(dst.MemoryUsageBytes, src.MemoryUsageBytes)
dst.CpuPercentMax = max(dst.CpuPercentMax, src.CpuPercentMax)
dst.MemoryPercentMax = max(dst.MemoryPercentMax, src.MemoryPercentMax)
dst.NetworkRxTotal = max(dst.NetworkRxTotal, src.NetworkRxTotal)
dst.NetworkTxTotal = max(dst.NetworkTxTotal, src.NetworkTxTotal)
dst.DiskReadTotal = max(dst.DiskReadTotal, src.DiskReadTotal)
dst.DiskWriteTotal = max(dst.DiskWriteTotal, src.DiskWriteTotal)
dst.CpuCores = max(dst.CpuCores, src.CpuCores)
dst.Samples += src.Samples
}
// startMetaRefresh rebuilds the metadata snapshot on its own goroutine and
// hands the result back over metaCh.
//
// This MUST stay off run()'s goroutine. ListAllContainers walks every local host
// sequentially with a 5s timeout each, so one slow or unreachable peer would
// stall the sample drain for N*5s. That stall does not stay local: the stats
// collector's dispatch loop blocks on `stats <- stat` for each subscriber in
// turn (internal/docker/stats_collector.go), so once our 512-deep sample buffer
// filled, a wedged cloud consumer would starve the LIVE UI stats broadcast on
// that host too. Metrics shipping must never be able to degrade the UI.
//
// At most one refresh is in flight; if the host service is slower than the
// window we keep serving the last-good snapshot rather than queueing work.
func (ss *statsStreamer) startMetaRefresh(ctx context.Context) {
if ss.refreshing {
return
}
ss.refreshing = true
go func() {
m := ss.buildMeta()
select {
case ss.metaCh <- m:
case <-ctx.Done():
}
}()
}
// refreshMeta rebuilds the snapshot synchronously. Only for tests and callers
// that are not on the sample-draining goroutine.
func (ss *statsStreamer) refreshMeta() {
ss.meta = ss.buildMeta()
}
// buildMeta resolves every container's name, CPU divisor and opt-out state.
// Called once per window rather than per sample — it is a host round-trip, and
// there are ~30 samples per container per window.
func (ss *statsStreamer) buildMeta() map[seriesKey]containerMeta {
containers, errs := ss.hostService.ListAllContainers(ss.labels)
for _, err := range errs {
if err != nil {
// Warned once, then quiet. A host that cannot be listed drops every
// one of its containers out of the batch two windows later (see
// flush) — silence that looks exactly like idle containers, which
// is precisely the confusion this whole retry pass exists to end.
// Repeating it every 30s for a permanently unreachable peer would
// bury everything else, so the rest go to Debug.
ss.listWarnOnce.Do(func() {
log.Warn().Err(err).Msg("stats streamer: error listing containers from host; its containers will not report metrics (further errors at debug level)")
})
log.Debug().Err(err).Msg("stats streamer: error listing containers from host")
}
}
ncpuByHost := make(map[string]int)
for _, h := range ss.hostService.Hosts() {
ncpuByHost[h.ID] = h.NCPU
}
meta := make(map[seriesKey]containerMeta, len(containers))
for _, c := range containers {
// A container the user opted out of log streaming with the min_level
// label shouldn't quietly keep reporting metrics either.
_, disabled, valid := parseMinLevel(c.Labels[cloudMinLevelLabel])
if !valid {
// logStreamer already logs the invalid value loudly for this
// container; don't double-report it every 30 seconds.
disabled = false
}
cores := c.CPULimit
if cores <= 0 {
cores = float64(ncpuByHost[c.Host])
}
if cores <= 0 {
cores = 1
}
meta[seriesKey{hostID: c.Host, containerID: c.ID}] = containerMeta{
name: c.Name,
cores: cores,
disabled: disabled,
}
}
return meta
}
+540
View File
@@ -0,0 +1,540 @@
package cloud
import (
"context"
"sync"
"sync/atomic"
"testing"
"time"
"github.com/amir20/dozzle/internal/container"
container_support "github.com/amir20/dozzle/internal/support/container"
pb "github.com/amir20/dozzle/proto/cloud"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
// fakeStatsHostService is a StatsStreamHostService returning a scripted
// container/host snapshot. SubscribeStats is unused by the tests below, which
// drive observe/flush directly rather than waiting on the 30s ticker.
type fakeStatsHostService struct {
containers []container.Container
hosts []container.Host
// subscribed, when non-nil, hands the test the channel run() is draining.
subscribed chan chan<- StatSample
}
func (f *fakeStatsHostService) ListAllContainers(_ container.ContainerLabels) ([]container.Container, []error) {
return f.containers, nil
}
func (f *fakeStatsHostService) FindContainer(_ string, _ string, _ container.ContainerLabels) (*container_support.ContainerService, error) {
return nil, nil
}
func (f *fakeStatsHostService) Hosts() []container.Host { return f.hosts }
func (f *fakeStatsHostService) SubscribeStats(_ context.Context, samples chan<- StatSample) {
if f.subscribed != nil {
f.subscribed <- samples
}
}
// newTestStreamer builds a streamer over one host with the given containers and
// a capture func standing in for the gRPC send.
func newTestStreamer(t *testing.T, containers []container.Container, ncpu int) (*statsStreamer, *[]*pb.StatsBatch) {
t.Helper()
hs := &fakeStatsHostService{
containers: containers,
hosts: []container.Host{{ID: "host1", Name: "host1", NCPU: ncpu}},
}
var sent []*pb.StatsBatch
ss := newStatsStreamer(hs, nil, func(resp *pb.ToolResponse) error {
sent = append(sent, resp.GetStatsBatch())
return nil
})
ss.refreshMeta()
return ss, &sent
}
func testContainer(id, name string) container.Container {
return container.Container{ID: id, Name: name, Host: "host1", Labels: map[string]string{}}
}
func TestStatsStreamer_AveragesGaugesAndKeepsLastCounters(t *testing.T) {
ss, sent := newTestStreamer(t, []container.Container{testContainer("c1", "api")}, 1)
for i, cpu := range []float64{10, 20, 60} {
ss.observe(StatSample{HostID: "host1", Stat: container.ContainerStat{
ID: "c1",
CPUPercent: cpu,
MemoryPercent: cpu / 2,
MemoryUsage: float64(100 * (i + 1)),
NetworkRxTotal: uint64(1000 * (i + 1)),
NetworkTxTotal: uint64(2000 * (i + 1)),
DiskReadTotal: uint64(10 * (i + 1)),
DiskWriteTotal: uint64(20 * (i + 1)),
}})
}
now := time.Unix(1700000000, 0)
require.NoError(t, ss.flush(now))
require.Len(t, *sent, 1)
entries := (*sent)[0].GetEntries()
require.Len(t, entries, 1)
e := entries[0]
assert.Equal(t, "host1", e.GetHostId())
assert.Equal(t, "api", e.GetContainerName())
assert.Equal(t, now.UnixNano(), e.GetTimestampNs())
assert.InDelta(t, 30.0, e.GetCpuPercent(), 0.001, "mean of 10/20/60")
assert.InDelta(t, 60.0, e.GetCpuPercentMax(), 0.001)
assert.InDelta(t, 15.0, e.GetMemoryPercent(), 0.001)
assert.InDelta(t, 30.0, e.GetMemoryPercentMax(), 0.001)
assert.InDelta(t, 200.0, e.GetMemoryUsageBytes(), 0.001, "mean of 100/200/300")
// Counters are cumulative — the newest value, not a sum or a delta.
assert.Equal(t, uint64(3000), e.GetNetworkRxTotal())
assert.Equal(t, uint64(6000), e.GetNetworkTxTotal())
assert.Equal(t, uint64(30), e.GetDiskReadTotal())
assert.Equal(t, uint64(60), e.GetDiskWriteTotal())
assert.Equal(t, uint32(3), e.GetSamples())
assert.InDelta(t, 1.0, e.GetCpuCores(), 0.001)
}
func TestStatsStreamer_NormalisesCPUByCores(t *testing.T) {
t.Run("falls back to host NCPU", func(t *testing.T) {
ss, sent := newTestStreamer(t, []container.Container{testContainer("c1", "api")}, 4)
ss.observe(StatSample{HostID: "host1", Stat: container.ContainerStat{ID: "c1", CPUPercent: 200}})
require.NoError(t, ss.flush(time.Unix(1, 0)))
e := (*sent)[0].GetEntries()[0]
assert.InDelta(t, 50.0, e.GetCpuPercent(), 0.001, "200%% of one core across 4 cores is 50%% overall")
assert.InDelta(t, 50.0, e.GetCpuPercentMax(), 0.001)
assert.InDelta(t, 4.0, e.GetCpuCores(), 0.001)
})
t.Run("prefers the container CPU limit", func(t *testing.T) {
c := testContainer("c1", "api")
c.CPULimit = 2
ss, sent := newTestStreamer(t, []container.Container{c}, 8)
ss.observe(StatSample{HostID: "host1", Stat: container.ContainerStat{ID: "c1", CPUPercent: 100}})
require.NoError(t, ss.flush(time.Unix(1, 0)))
e := (*sent)[0].GetEntries()[0]
assert.InDelta(t, 50.0, e.GetCpuPercent(), 0.001)
assert.InDelta(t, 2.0, e.GetCpuCores(), 0.001)
})
t.Run("defaults to one core when nothing is known", func(t *testing.T) {
ss, sent := newTestStreamer(t, []container.Container{testContainer("c1", "api")}, 0)
ss.observe(StatSample{HostID: "host1", Stat: container.ContainerStat{ID: "c1", CPUPercent: 37}})
require.NoError(t, ss.flush(time.Unix(1, 0)))
assert.InDelta(t, 37.0, (*sent)[0].GetEntries()[0].GetCpuPercent(), 0.001)
})
}
func TestStatsStreamer_ResetsBetweenWindows(t *testing.T) {
ss, sent := newTestStreamer(t, []container.Container{testContainer("c1", "api")}, 1)
ss.observe(StatSample{HostID: "host1", Stat: container.ContainerStat{ID: "c1", CPUPercent: 80}})
require.NoError(t, ss.flush(time.Unix(1, 0)))
assert.InDelta(t, 80.0, (*sent)[0].GetEntries()[0].GetCpuPercent(), 0.001)
// Second window sees a single low sample — if state leaked the mean would
// still be dragged up by the previous window's 80.
ss.observe(StatSample{HostID: "host1", Stat: container.ContainerStat{ID: "c1", CPUPercent: 10}})
require.NoError(t, ss.flush(time.Unix(2, 0)))
require.Len(t, *sent, 2)
assert.InDelta(t, 10.0, (*sent)[1].GetEntries()[0].GetCpuPercent(), 0.001)
}
func TestStatsStreamer_EmptyWindowSendsNothing(t *testing.T) {
ss, sent := newTestStreamer(t, []container.Container{testContainer("c1", "api")}, 1)
require.NoError(t, ss.flush(time.Unix(1, 0)))
assert.Empty(t, *sent, "a window with no samples must not push an empty batch")
// A container that stops simply stops producing samples, so it drops out.
ss.observe(StatSample{HostID: "host1", Stat: container.ContainerStat{ID: "c1", CPUPercent: 5}})
require.NoError(t, ss.flush(time.Unix(2, 0)))
require.Len(t, *sent, 1)
require.NoError(t, ss.flush(time.Unix(3, 0)))
assert.Len(t, *sent, 1, "no new batch after the container went quiet")
}
func TestStatsStreamer_SkipsDisabledContainers(t *testing.T) {
disabled := testContainer("c1", "secret")
disabled.Labels[cloudMinLevelLabel] = "disabled"
ss, sent := newTestStreamer(t, []container.Container{disabled, testContainer("c2", "api")}, 1)
ss.observe(StatSample{HostID: "host1", Stat: container.ContainerStat{ID: "c1", CPUPercent: 90}})
ss.observe(StatSample{HostID: "host1", Stat: container.ContainerStat{ID: "c2", CPUPercent: 10}})
require.NoError(t, ss.flush(time.Unix(1, 0)))
require.Len(t, *sent, 1)
entries := (*sent)[0].GetEntries()
require.Len(t, entries, 1)
assert.Equal(t, "api", entries[0].GetContainerName())
}
func TestStatsStreamer_HoldsThenDropsUnresolvedContainers(t *testing.T) {
// No containers in the snapshot at all — nothing can resolve to a name.
ss, sent := newTestStreamer(t, nil, 1)
ss.observe(StatSample{HostID: "host1", Stat: container.ContainerStat{ID: "ghost", CPUPercent: 50}})
require.NoError(t, ss.flush(time.Unix(1, 0)))
assert.Empty(t, *sent)
assert.Len(t, ss.acc, 1, "held for one grace window in case metadata catches up")
require.NoError(t, ss.flush(time.Unix(2, 0)))
assert.Empty(t, *sent)
assert.Empty(t, ss.acc, "abandoned after the grace window rather than leaking")
}
func TestStatsStreamer_EmitsOnceMetadataCatchesUp(t *testing.T) {
hs := &fakeStatsHostService{hosts: []container.Host{{ID: "host1", NCPU: 1}}}
var sent []*pb.StatsBatch
ss := newStatsStreamer(hs, nil, func(resp *pb.ToolResponse) error {
sent = append(sent, resp.GetStatsBatch())
return nil
})
ss.refreshMeta()
// Container started mid-window: samples arrive before it shows up in a listing.
ss.observe(StatSample{HostID: "host1", Stat: container.ContainerStat{ID: "c1", CPUPercent: 40}})
require.NoError(t, ss.flush(time.Unix(1, 0)))
assert.Empty(t, sent)
hs.containers = []container.Container{testContainer("c1", "late-starter")}
ss.refreshMeta() // stands in for the async refresh run() kicks off each tick
ss.observe(StatSample{HostID: "host1", Stat: container.ContainerStat{ID: "c1", CPUPercent: 60}})
require.NoError(t, ss.flush(time.Unix(2, 0)))
require.Len(t, sent, 1)
e := sent[0].GetEntries()[0]
assert.Equal(t, "late-starter", e.GetContainerName())
assert.InDelta(t, 50.0, e.GetCpuPercent(), 0.001, "both windows' samples are folded in")
assert.Equal(t, uint32(2), e.GetSamples())
}
func TestStatsStreamer_SeparatesSameNameOnDifferentHosts(t *testing.T) {
hs := &fakeStatsHostService{
containers: []container.Container{
{ID: "c1", Name: "api", Host: "host1", Labels: map[string]string{}},
{ID: "c2", Name: "api", Host: "host2", Labels: map[string]string{}},
},
hosts: []container.Host{{ID: "host1", NCPU: 1}, {ID: "host2", NCPU: 1}},
}
var sent []*pb.StatsBatch
ss := newStatsStreamer(hs, nil, func(resp *pb.ToolResponse) error {
sent = append(sent, resp.GetStatsBatch())
return nil
})
ss.refreshMeta()
ss.observe(StatSample{HostID: "host1", Stat: container.ContainerStat{ID: "c1", CPUPercent: 10}})
ss.observe(StatSample{HostID: "host2", Stat: container.ContainerStat{ID: "c2", CPUPercent: 90}})
require.NoError(t, ss.flush(time.Unix(1, 0)))
entries := sent[0].GetEntries()
require.Len(t, entries, 2)
// Sorted by (host_id, container_name).
assert.Equal(t, "host1", entries[0].GetHostId())
assert.InDelta(t, 10.0, entries[0].GetCpuPercent(), 0.001)
assert.Equal(t, "host2", entries[1].GetHostId())
assert.InDelta(t, 90.0, entries[1].GetCpuPercent(), 0.001)
}
// A container redeployed inside one window keeps its name and gets a new ID, so
// the dying and the starting container hold two accumulators that resolve to the
// same series. They have to collapse into one entry — two points on the same
// (host, name) series at the same timestamp is not something the wire format can
// express, since StatsBatchEntry carries no container ID.
func TestStatsStreamer_MergesRedeployedContainerWithinAWindow(t *testing.T) {
ss, sent := newTestStreamer(t, []container.Container{
testContainer("old", "api"),
testContainer("new", "api"),
}, 1)
// Old container: 3 samples, high counters from a long uptime.
for _, cpu := range []float64{10, 20, 30} {
ss.observe(StatSample{HostID: "host1", Stat: container.ContainerStat{
ID: "old", CPUPercent: cpu, MemoryPercent: 40, MemoryUsage: 100,
NetworkRxTotal: 9000, DiskReadTotal: 900,
}})
}
// Replacement: 1 sample, counters back near zero.
ss.observe(StatSample{HostID: "host1", Stat: container.ContainerStat{
ID: "new", CPUPercent: 80, MemoryPercent: 80, MemoryUsage: 300,
NetworkRxTotal: 5, DiskReadTotal: 1,
}})
require.NoError(t, ss.flush(time.Unix(1, 0)))
entries := (*sent)[0].GetEntries()
require.Len(t, entries, 1, "one series, not two points at the same timestamp")
e := entries[0]
assert.Equal(t, "api", e.GetContainerName())
assert.Equal(t, uint32(4), e.GetSamples())
// Means weighted by sample count, not a flat average of the two entries.
assert.InDelta(t, 35.0, e.GetCpuPercent(), 0.001, "(10+20+30+80)/4")
assert.InDelta(t, 50.0, e.GetMemoryPercent(), 0.001, "(40*3+80)/4")
assert.InDelta(t, 150.0, e.GetMemoryUsageBytes(), 0.001, "(100*3+300)/4")
assert.InDelta(t, 80.0, e.GetCpuPercentMax(), 0.001)
assert.InDelta(t, 80.0, e.GetMemoryPercentMax(), 0.001)
// Counters take the larger: the series was already sitting at the old
// container's total, and summing would report a value neither ever had.
assert.Equal(t, uint64(9000), e.GetNetworkRxTotal())
assert.Equal(t, uint64(900), e.GetDiskReadTotal())
}
// slowStatsHostService lets the first N ListAllContainers calls through, then
// blocks — standing in for a multi-host setup where one peer goes unreachable
// mid-run (each host costs a 5s timeout, sequentially).
type slowStatsHostService struct {
fakeStatsHostService
blockAfter int32
calls atomic.Int32
entered chan struct{}
release chan struct{}
once sync.Once
}
func (f *slowStatsHostService) ListAllContainers(l container.ContainerLabels) ([]container.Container, []error) {
if f.calls.Add(1) > f.blockAfter {
f.once.Do(func() { close(f.entered) })
<-f.release
}
return f.fakeStatsHostService.ListAllContainers(l)
}
// A slow host must not stall the sample-draining loop.
//
// flush() used to call ListAllContainers inline, on the same goroutine that
// drains samples. A hung host would block that drain; once the sample buffer
// filled, backpressure reached the stats collector's blocking per-subscriber
// dispatch and starved the LIVE UI stats broadcast on that host. Metrics
// shipping must never be able to degrade the UI.
//
// The window is driven short here so flush() actually runs — that is where the
// original bug lived.
func TestStatsStreamer_SlowHostDoesNotStallSampleDrain(t *testing.T) {
hs := &slowStatsHostService{
fakeStatsHostService: fakeStatsHostService{
containers: []container.Container{testContainer("c1", "api")},
hosts: []container.Host{{ID: "host1", NCPU: 1}},
subscribed: make(chan chan<- StatSample, 1),
},
blockAfter: 1, // let the startup refresh through, wedge the flush-time one
entered: make(chan struct{}),
release: make(chan struct{}),
}
ss := newStatsStreamer(hs, nil, func(*pb.ToolResponse) error { return nil })
ss.window = 20 * time.Millisecond
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
done := make(chan struct{})
go func() { defer close(done); ss.run(ctx) }()
samples := <-hs.subscribed
// Wait for a refresh to wedge, which only happens once the ticker has fired.
select {
case <-hs.entered:
case <-time.After(5 * time.Second):
t.Fatal("host never blocked; test did not reach the interesting state")
}
// With a refresh wedged, the loop must still drain. Push well past the
// buffer: if the drain were blocked behind the host call this never finishes.
drained := make(chan struct{})
go func() {
defer close(drained)
for range statsSampleChanBuf * 3 {
select {
case samples <- StatSample{HostID: "host1", Stat: container.ContainerStat{ID: "c1", CPUPercent: 1}}:
case <-ctx.Done():
return
}
}
}()
select {
case <-drained:
case <-time.After(5 * time.Second):
t.Fatal("sample drain stalled behind a slow host — backpressure would reach the stats collector and starve the live UI stats broadcast")
}
close(hs.release)
cancel()
<-done
}
// The refreshed snapshot must actually be installed by the run loop — if it
// weren't, names would never resolve and every window would be dropped. Asserted
// through run()'s real ticker on observable output, rather than by reading
// ss.meta, which the run goroutine owns.
func TestStatsStreamer_AsyncRefreshInstallsSnapshot(t *testing.T) {
hs := &fakeStatsHostService{
containers: []container.Container{testContainer("c1", "api")},
hosts: []container.Host{{ID: "host1", NCPU: 1}},
subscribed: make(chan chan<- StatSample, 1),
}
batches := make(chan *pb.StatsBatch, 4)
ss := newStatsStreamer(hs, nil, func(r *pb.ToolResponse) error {
batches <- r.GetStatsBatch()
return nil
})
ss.window = 50 * time.Millisecond
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
done := make(chan struct{})
go func() { defer close(done); ss.run(ctx) }()
samples := <-hs.subscribed
samples <- StatSample{HostID: "host1", Stat: container.ContainerStat{ID: "c1", CPUPercent: 42}}
select {
case b := <-batches:
require.Len(t, b.GetEntries(), 1)
assert.Equal(t, "api", b.GetEntries()[0].GetContainerName(),
"name resolved, so the async snapshot was installed")
assert.InDelta(t, 42.0, b.GetEntries()[0].GetCpuPercent(), 0.001)
case <-time.After(5 * time.Second):
t.Fatal("no batch emitted — async refresh never installed its snapshot")
}
cancel()
<-done
}
// A refresh slower than the window must not queue up more refreshes.
func TestStatsStreamer_OnlyOneRefreshInFlight(t *testing.T) {
hs := &slowStatsHostService{
fakeStatsHostService: fakeStatsHostService{
hosts: []container.Host{{ID: "host1", NCPU: 1}},
},
blockAfter: 0, // block immediately
entered: make(chan struct{}),
release: make(chan struct{}),
}
ss := newStatsStreamer(hs, nil, func(*pb.ToolResponse) error { return nil })
ctx := t.Context()
ss.startMetaRefresh(ctx)
<-hs.entered
// Every subsequent attempt while one is in flight is a no-op.
for range 5 {
ss.startMetaRefresh(ctx)
}
assert.True(t, ss.refreshing, "should still be marked in-flight")
close(hs.release)
select {
case m := <-ss.metaCh:
assert.NotNil(t, m, "exactly one snapshot delivered")
case <-time.After(2 * time.Second):
t.Fatal("refresh never completed")
}
// Nothing else queued behind it.
select {
case <-ss.metaCh:
t.Fatal("a second refresh was queued while one was in flight")
case <-time.After(100 * time.Millisecond):
}
}
// A failed send must not end the streamer. It used to `return`, so a single
// transient error meant the instance reported no metrics at all until it
// reconnected — the failure mode that made a live Cloud instance look idle for
// hours.
func TestStatsStreamer_SurvivesAFailedSend(t *testing.T) {
hs := &fakeStatsHostService{
containers: []container.Container{testContainer("c1", "api")},
hosts: []container.Host{{ID: "host1", NCPU: 1}},
subscribed: make(chan chan<- StatSample, 1),
}
var failNext atomic.Bool
failNext.Store(true)
batches := make(chan *pb.StatsBatch, 4)
ss := newStatsStreamer(hs, nil, func(r *pb.ToolResponse) error {
if failNext.Swap(false) {
return assert.AnError
}
batches <- r.GetStatsBatch()
return nil
})
ss.window = 50 * time.Millisecond
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
done := make(chan struct{})
go func() { defer close(done); ss.run(ctx) }()
samples := <-hs.subscribed
// Two windows' worth: the first send fails, the second must still happen.
samples <- StatSample{HostID: "host1", Stat: container.ContainerStat{ID: "c1", CPUPercent: 10}}
time.Sleep(80 * time.Millisecond)
samples <- StatSample{HostID: "host1", Stat: container.ContainerStat{ID: "c1", CPUPercent: 20}}
select {
case b := <-batches:
require.Len(t, b.GetEntries(), 1)
assert.Equal(t, "api", b.GetEntries()[0].GetContainerName())
case <-time.After(5 * time.Second):
t.Fatal("streamer stopped after one failed send")
}
cancel()
<-done
}
// ...but a stream that fails every single time is wedged, and spinning on it
// forever helps nobody. It gives up after statsMaxSendFailures.
func TestStatsStreamer_GivesUpOnAPermanentlyFailingSend(t *testing.T) {
hs := &fakeStatsHostService{
containers: []container.Container{testContainer("c1", "api")},
hosts: []container.Host{{ID: "host1", NCPU: 1}},
subscribed: make(chan chan<- StatSample, 1),
}
var sends atomic.Int32
ss := newStatsStreamer(hs, nil, func(*pb.ToolResponse) error {
sends.Add(1)
return assert.AnError
})
ss.window = 20 * time.Millisecond
ctx := t.Context()
done := make(chan struct{})
go func() { defer close(done); ss.run(ctx) }()
samples := <-hs.subscribed
// Keep feeding so every window has something to send; run() exits once the
// failures reach the cap.
go func() {
for {
select {
case <-ctx.Done():
return
case samples <- StatSample{HostID: "host1", Stat: container.ContainerStat{ID: "c1", CPUPercent: 1}}:
time.Sleep(5 * time.Millisecond)
}
}
}()
select {
case <-done:
assert.Equal(t, int32(statsMaxSendFailures), sends.Load(),
"should stop exactly at the failure cap")
case <-time.After(5 * time.Second):
t.Fatal("streamer never gave up on a permanently failing send")
}
}
+104 -10
View File
@@ -27,6 +27,37 @@ type DockerStatsCollector struct {
var timeToStop = 6 * time.Hour
// Retry bounds for the two Docker streams this collector depends on: a
// container's stats stream and the host's event stream.
//
// Both used to be one-shot. That was survivable when the only subscribers were
// UI tabs — a user looking at a dead chart reloads the page, which resubscribes
// and rebuilds everything. It is not survivable now that Cloud metrics keep a
// permanent subscription: nothing ever resubscribes, so a single transient
// error meant that container (or the whole host) stopped reporting until the
// process restarted, with no signal anywhere that it had.
var (
streamRetryMin = 1 * time.Second
streamRetryMax = 30 * time.Second
)
// nextBackoff doubles up to the ceiling.
func nextBackoff(d time.Duration) time.Duration {
return min(d*2, streamRetryMax)
}
// sleepOrDone waits out the backoff, returning false if ctx ended first.
func sleepOrDone(ctx context.Context, d time.Duration) bool {
t := time.NewTimer(d)
defer t.Stop()
select {
case <-ctx.Done():
return false
case <-t.C:
return true
}
}
func NewDockerStatsCollector(client container.Client, labels container.ContainerLabels) *DockerStatsCollector {
return &DockerStatsCollector{
stream: make(chan container.ContainerStat),
@@ -74,15 +105,47 @@ func (c *DockerStatsCollector) reset() {
c.timer = nil
}
// streamStats keeps one container's stats flowing for as long as its context
// lives. ContainerStats returning — with an error OR cleanly — means the stream
// broke, not that the container is gone: a container that actually stops
// arrives as a `die` event, which cancels this context. So anything else is
// retried with backoff. Without that, one hiccup left a single container
// silently absent from every subsequent stats window while its neighbours
// carried on, which is indistinguishable from a container that is simply idle.
func streamStats(parent context.Context, sc *DockerStatsCollector, id string) {
ctx, cancel := context.WithCancel(parent)
sc.cancelers.Store(id, cancel)
log.Debug().Str("container", id).Str("host", sc.client.Host().Name).Msg("starting to stream stats")
if err := sc.client.ContainerStats(ctx, id, sc.stream); err != nil {
log.Debug().Str("container", id).Str("host", sc.client.Host().Name).Err(err).Msg("stopping to stream stats")
if !errors.Is(err, context.Canceled) && !errors.Is(err, io.EOF) {
log.Error().Str("container", id).Str("host", sc.client.Host().Name).Err(err).Msg("unexpected error while streaming stats")
backoff := streamRetryMin
for {
log.Debug().Str("container", id).Str("host", sc.client.Host().Name).Msg("starting to stream stats")
err := sc.client.ContainerStats(ctx, id, sc.stream)
// Cancelled = the container died or the collector shut down. Expected,
// and the only way out of this loop.
if ctx.Err() != nil {
log.Debug().Str("container", id).Str("host", sc.client.Host().Name).Msg("stopping to stream stats")
return
}
// Warn, not Debug: this is the failure that used to be permanent, and
// the whole point of the retry is that someone can see it happening.
ev := log.Warn()
if err == nil || errors.Is(err, io.EOF) {
// A clean end is ordinary (Docker closes idle stats streams), so it
// does not deserve a warning until it starts repeating.
ev = log.Debug()
}
ev.Str("container", id).
Str("host", sc.client.Host().Name).
Err(err).
Dur("retry_in", backoff).
Msg("container stats stream ended, retrying")
if !sleepOrDone(ctx, backoff) {
return
}
backoff = nextBackoff(backoff)
}
}
@@ -114,20 +177,51 @@ func (sc *DockerStatsCollector) Start(parentCtx context.Context) bool {
events := make(chan container.ContainerEvent)
// The event stream is how new containers get a stats stream and dead ones
// lose theirs, so losing it degrades the collector even while the existing
// per-container streams keep running. It used to forceStop() the whole
// collector on any error, which bypasses the reference count entirely:
// every subscriber — including a Cloud connection that is never coming back
// to resubscribe — went silent at once and was never told. Now it
// reconnects, and only a cancelled context ends it.
go func() {
defer close(events)
log.Debug().Str("host", sc.client.Host().Name).Msg("starting to listen to docker events")
err := sc.client.ContainerEvents(ctx, events)
if err != nil && !errors.Is(err, context.Canceled) {
log.Error().Str("host", sc.client.Host().Name).Err(err).Msg("unexpected error while listening to docker events")
backoff := streamRetryMin
for {
log.Debug().Str("host", sc.client.Host().Name).Msg("starting to listen to docker events")
err := sc.client.ContainerEvents(ctx, events)
if ctx.Err() != nil {
return
}
// A clean exit is ordinary, so it only gets a warning once it starts
// repeating for a real reason.
ev := log.Warn()
if err == nil {
ev = log.Debug()
}
ev.
Str("host", sc.client.Host().Name).
Err(err).
Dur("retry_in", backoff).
Msg("docker event stream ended, retrying")
if !sleepOrDone(ctx, backoff) {
return
}
backoff = nextBackoff(backoff)
}
sc.forceStop()
}()
go func() {
for event := range events {
switch event.Name {
case "start":
// Replace any stream still retrying for this id. Without this,
// a start that arrives without a matching die (a restart, a
// duplicate event) would leave two loops feeding sc.stream for
// one container, doubling its samples.
if cancel, ok := sc.cancelers.LoadAndDelete(event.ActorID); ok {
cancel()
}
go streamStats(ctx, sc, event.ActorID)
case "die":
+101
View File
@@ -2,7 +2,9 @@ package docker
import (
"context"
"errors"
"testing"
"time"
"github.com/amir20/dozzle/internal/container"
"github.com/stretchr/testify/assert"
@@ -113,3 +115,102 @@ func TestStop(t *testing.T) {
collector.Stop()
assert.Equal(t, int32(0), collector.totalStarted.Load(), "total started should be 1")
}
// fastRetries shrinks the backoff so a retry test finishes in milliseconds.
func fastRetries(t *testing.T) {
t.Helper()
oldMin, oldMax := streamRetryMin, streamRetryMax
streamRetryMin, streamRetryMax = time.Millisecond, time.Millisecond
t.Cleanup(func() { streamRetryMin, streamRetryMax = oldMin, oldMax })
}
// A container's stats stream breaking is not the container going away — that
// arrives as a `die` event. This used to be one-shot, so a single transient
// error left one container permanently absent from stats while every other
// container on the host carried on: indistinguishable from an idle container,
// and unrecoverable without restarting the process.
func TestStatsStreamRetriesAfterError(t *testing.T) {
fastRetries(t)
ctx, cancel := context.WithCancel(context.Background())
t.Cleanup(cancel)
client := new(mockedClient)
client.On("Host").Return(container.Host{ID: "localhost"})
client.On("ListContainers", mock.Anything, mock.Anything).Return([]container.Container{
{ID: "1234", Name: "test", State: "running"},
}, nil)
client.On("ContainerEvents", mock.Anything, mock.Anything).
Return(nil).
Run(func(args mock.Arguments) { <-args.Get(0).(context.Context).Done() })
// First attempt fails, every later attempt delivers.
client.On("ContainerStats", mock.Anything, mock.Anything, mock.Anything).
Return(errors.New("connection reset")).
Once()
client.On("ContainerStats", mock.Anything, mock.Anything, mock.Anything).
Return(nil).
Run(func(args mock.Arguments) {
args.Get(2).(chan<- container.ContainerStat) <- container.ContainerStat{ID: "1234"}
})
collector := NewDockerStatsCollector(client, container.ContainerLabels{})
stats := make(chan container.ContainerStat)
collector.Subscribe(ctx, stats)
go collector.Start(ctx)
select {
case s := <-stats:
assert.Equal(t, "1234", s.ID, "stats resumed after the stream errored")
case <-time.After(5 * time.Second):
t.Fatal("stats never resumed after a failed stream — the retry is gone")
}
}
// The event stream ending used to call forceStop(), which bypasses the
// reference count: every subscriber went silent at once and none was told. A
// Cloud connection holds a permanent subscription and never resubscribes, so
// that was terminal for it. It must reconnect instead.
func TestEventStreamRetriesInsteadOfStoppingTheCollector(t *testing.T) {
fastRetries(t)
ctx, cancel := context.WithCancel(context.Background())
t.Cleanup(cancel)
reconnected := make(chan struct{}, 1)
client := new(mockedClient)
client.On("Host").Return(container.Host{ID: "localhost"})
client.On("ListContainers", mock.Anything, mock.Anything).Return([]container.Container{}, nil)
client.On("ContainerStats", mock.Anything, mock.Anything, mock.Anything).Return(nil).Maybe()
client.On("ContainerEvents", mock.Anything, mock.Anything).
Return(errors.New("docker went away")).
Once()
client.On("ContainerEvents", mock.Anything, mock.Anything).
Return(nil).
Run(func(args mock.Arguments) {
select {
case reconnected <- struct{}{}:
default:
}
<-args.Get(0).(context.Context).Done()
})
collector := NewDockerStatsCollector(client, container.ContainerLabels{})
stopped := make(chan bool, 1)
go func() { stopped <- collector.Start(ctx) }()
select {
case <-reconnected:
case <-time.After(5 * time.Second):
t.Fatal("event stream was never re-established")
}
// And the collector itself is still up — the old behaviour tore it down.
select {
case <-stopped:
t.Fatal("collector stopped when the event stream failed")
case <-time.After(50 * time.Millisecond):
}
collector.mu.Lock()
assert.NotNil(t, collector.stopper, "collector should still be running")
collector.mu.Unlock()
}
+10 -2
View File
@@ -358,8 +358,16 @@ cloud:
error-unavailable: Dozzle Cloud er midlertidigt utilgængelig. Prøv igen senere.
unlink: Fjern tilknytning
unlink-confirm: Er du sikker på, at du vil fjerne tilknytningen til Dozzle Cloud? Dette fjerner alle cloud-notifikationsdestinationer.
stream-logs: Stream containerlogs til Dozzle Cloud
stream-logs-help: Påkrævet for AI-drevne undersøgelser og logsøgning. Deaktiver for at holde alt logindhold på denne instans.
privacy:
toggle: "Stream logs og metrikker til Dozzle Cloud"
sends-heading: "Så længe dette er slået til, sender dine containere løbende:"
sends-logs: "Loglinjer, efterhånden som de skrives"
sends-metrics: "CPU-, hukommelses-, netværks- og diskforbrug, gennemsnitligt over 30 sekunder"
never-heading: "Sendes aldrig:"
never-env: "Miljøvariabler"
off-note: "Slå dette fra for at holde logs og metrikker på denne instans. Alarmer, du har opsat, leveres stadig, og fjernadgang til containere fungerer fortsat."
per-container: "Udelad en enkelt container med labelet: dev.dozzle.cloud.min_level=disabled"
required-for: "Påkrævet for AI-undersøgelser, logsøgning og ressourcehistorik."
welcome:
title: "Cloud er en intelligent triage for dine containere."
subtitle: "Et startsæt af fejlsignaler er forhåndsvalgt — fortæl os, hvad der bragte dig hertil, så vi kan finjustere det."
+10 -2
View File
@@ -358,8 +358,16 @@ cloud:
error-unavailable: Dozzle Cloud ist vorübergehend nicht verfügbar. Bitte versuche es später erneut.
unlink: Verknüpfung aufheben
unlink-confirm: Sind Sie sicher, dass Sie die Verknüpfung mit Dozzle Cloud aufheben möchten? Dies entfernt alle Cloud-Benachrichtigungsziele.
stream-logs: Container-Logs an Dozzle Cloud streamen
stream-logs-help: Erforderlich für KI-gestützte Analysen und Log-Suche. Deaktivieren, um sämtliche Log-Inhalte auf dieser Instanz zu behalten.
privacy:
toggle: "Logs und Metriken an Dozzle Cloud streamen"
sends-heading: "Solange dies aktiviert ist, senden Ihre Container fortlaufend:"
sends-logs: "Log-Zeilen, sobald sie geschrieben werden"
sends-metrics: "CPU-, Arbeitsspeicher-, Netzwerk- und Festplattennutzung, gemittelt über 30 Sekunden"
never-heading: "Wird nie gesendet:"
never-env: "Umgebungsvariablen"
off-note: "Deaktivieren Sie dies, um Logs und Metriken auf dieser Instanz zu behalten. Eingerichtete Benachrichtigungen werden weiterhin zugestellt und der Fernzugriff auf Container funktioniert weiterhin."
per-container: "Einen einzelnen Container ausschließen mit dem Label: dev.dozzle.cloud.min_level=disabled"
required-for: "Erforderlich für KI-Analysen, Log-Suche und Ressourcenverlauf."
welcome:
title: "Cloud ist eine intelligente Triage für Ihre Container."
subtitle: "Ein Starter-Set an Fehlersignalen ist vorausgewählt — sagen Sie uns, was Sie hergeführt hat, damit wir es feinabstimmen können."
+10 -2
View File
@@ -409,8 +409,16 @@ cloud:
error-unavailable: Dozzle Cloud is temporarily unavailable. Please try again later.
unlink: Unlink
unlink-confirm: Are you sure you want to unlink from Dozzle Cloud? This will remove all cloud notification destinations.
stream-logs: Stream container logs to Dozzle Cloud
stream-logs-help: Required for AI-powered investigations and log search. Disable to keep all log content on this instance.
privacy:
toggle: Stream logs and metrics to Dozzle Cloud
sends-heading: "While this is on, your containers continuously send:"
sends-logs: Log lines, as they are written
sends-metrics: CPU, memory, network and disk usage, averaged over 30 seconds
never-heading: "Never sent:"
never-env: Environment variables
off-note: Turn this off to keep logs and metrics on this instance. Alerts you have set up are still delivered, and remote container access still works.
per-container: "Exclude a single container with the label: dev.dozzle.cloud.min_level=disabled"
required-for: Required for AI investigations, log search, and resource history.
welcome:
title: "Cloud is an intelligent triage for your containers."
subtitle: "A starter set of failure signals is pre-selected — tell us what brought you here so we can tune it."
+10 -2
View File
@@ -402,8 +402,16 @@ cloud:
error-unavailable: Dozzle Cloud no está disponible temporalmente. Inténtalo de nuevo más tarde.
unlink: Desvincular
unlink-confirm: ¿Está seguro de que desea desvincular de Dozzle Cloud? Esto eliminará todos los destinos de notificación en la nube.
stream-logs: Transmitir registros de contenedores a Dozzle Cloud
stream-logs-help: Necesario para investigaciones con IA y búsqueda de registros. Desactívelo para mantener todo el contenido de los registros en esta instancia.
privacy:
toggle: "Transmitir registros y métricas a Dozzle Cloud"
sends-heading: "Mientras esto esté activado, sus contenedores envían continuamente:"
sends-logs: "Líneas de registro, a medida que se escriben"
sends-metrics: "Uso de CPU, memoria, red y disco, promediado en 30 segundos"
never-heading: "Nunca se envía:"
never-env: "Variables de entorno"
off-note: "Desactívelo para mantener los registros y las métricas en esta instancia. Las alertas que haya configurado se siguen entregando y el acceso remoto a los contenedores sigue funcionando."
per-container: "Excluya un contenedor concreto con la etiqueta: dev.dozzle.cloud.min_level=disabled"
required-for: "Necesario para investigaciones con IA, búsqueda de registros e historial de recursos."
welcome:
title: "Cloud es un triaje inteligente para tus contenedores."
subtitle: "Un conjunto inicial de señales de fallo está preseleccionado — cuéntanos qué te trajo aquí para que podamos ajustarlo."
+10 -2
View File
@@ -358,8 +358,16 @@ cloud:
error-unavailable: Dozzle Cloud est temporairement indisponible. Veuillez réessayer plus tard.
unlink: Délier
unlink-confirm: Êtes-vous sûr de vouloir délier de Dozzle Cloud ? Cela supprimera toutes les destinations de notification cloud.
stream-logs: Diffuser les journaux de conteneurs vers Dozzle Cloud
stream-logs-help: Requis pour les investigations alimentées par IA et la recherche dans les journaux. Désactivez pour conserver tout le contenu des journaux sur cette instance.
privacy:
toggle: "Diffuser les journaux et les métriques vers Dozzle Cloud"
sends-heading: "Tant que cette option est activée, vos conteneurs envoient en continu :"
sends-logs: "Les lignes de journal, au fur et à mesure de leur écriture"
sends-metrics: "L'utilisation du processeur, de la mémoire, du réseau et du disque, moyennée sur 30 secondes"
never-heading: "Jamais envoyé :"
never-env: "Les variables d'environnement"
off-note: "Désactivez pour conserver les journaux et les métriques sur cette instance. Les alertes que vous avez configurées continuent d'être envoyées et l'accès distant aux conteneurs fonctionne toujours."
per-container: "Exclure un seul conteneur avec le label : dev.dozzle.cloud.min_level=disabled"
required-for: "Requis pour les investigations par IA, la recherche dans les journaux et l'historique des ressources."
welcome:
title: "Cloud est un triage intelligent pour vos conteneurs."
subtitle: "Un ensemble initial de signaux de défaillance est présélectionné — dites-nous ce qui vous amène pour que nous puissions l'ajuster."
+10 -2
View File
@@ -370,8 +370,16 @@ cloud:
error-unavailable: Dozzle Cloud sedang tidak tersedia. Silakan coba lagi nanti.
unlink: Lepaskan tautan
unlink-confirm: Apakah Anda yakin ingin melepaskan tautan dari Dozzle Cloud? Ini akan menghapus semua tujuan notifikasi cloud.
stream-logs: Streaming log kontainer ke Dozzle Cloud
stream-logs-help: Diperlukan untuk investigasi bertenaga AI dan pencarian log. Nonaktifkan untuk menyimpan seluruh konten log di instance ini.
privacy:
toggle: "Streaming log dan metrik ke Dozzle Cloud"
sends-heading: "Selama ini aktif, kontainer Anda terus mengirim:"
sends-logs: "Baris log, saat ditulis"
sends-metrics: "Penggunaan CPU, memori, jaringan, dan disk, dirata-ratakan selama 30 detik"
never-heading: "Tidak pernah dikirim:"
never-env: "Variabel lingkungan"
off-note: "Nonaktifkan untuk menyimpan log dan metrik di instance ini. Peringatan yang sudah Anda atur tetap dikirim, dan akses kontainer jarak jauh tetap berfungsi."
per-container: "Kecualikan satu kontainer dengan label: dev.dozzle.cloud.min_level=disabled"
required-for: "Diperlukan untuk investigasi AI, pencarian log, dan riwayat sumber daya."
welcome:
title: "Cloud adalah triage cerdas untuk kontainer Anda."
subtitle: "Sekumpulan sinyal kegagalan awal sudah dipilih — beri tahu kami apa yang membawa Anda kemari agar bisa kami sesuaikan."
+10 -2
View File
@@ -358,8 +358,16 @@ cloud:
error-unavailable: Dozzle Cloud non è temporaneamente disponibile. Riprova più tardi.
unlink: Scollega
unlink-confirm: Sei sicuro di voler scollegare da Dozzle Cloud? Questo rimuoverà tutte le destinazioni di notifica cloud.
stream-logs: Invia i log dei container a Dozzle Cloud
stream-logs-help: Necessario per le indagini basate sull'IA e la ricerca nei log. Disattiva per mantenere tutto il contenuto dei log su questa istanza.
privacy:
toggle: "Invia log e metriche a Dozzle Cloud"
sends-heading: "Finché questa opzione è attiva, i tuoi container inviano continuamente:"
sends-logs: "Le righe di log, man mano che vengono scritte"
sends-metrics: "Utilizzo di CPU, memoria, rete e disco, mediato su 30 secondi"
never-heading: "Non viene mai inviato:"
never-env: "Le variabili d'ambiente"
off-note: "Disattiva per mantenere log e metriche su questa istanza. Gli avvisi che hai configurato continuano a essere recapitati e l'accesso remoto ai container funziona ancora."
per-container: "Escludi un singolo container con l'etichetta: dev.dozzle.cloud.min_level=disabled"
required-for: "Necessario per le indagini con IA, la ricerca nei log e lo storico delle risorse."
welcome:
title: "Cloud è un triage intelligente per i tuoi container."
subtitle: "Un set iniziale di segnali di errore è preselezionato — dicci cosa ti ha portato qui così possiamo affinarlo."
+10 -2
View File
@@ -362,8 +362,16 @@ cloud:
error-unavailable: Dozzle Cloud를 일시적으로 사용할 수 없습니다. 나중에 다시 시도해 주세요.
unlink: 연결 해제
unlink-confirm: Dozzle Cloud 연결을 해제하시겠습니까? 모든 클라우드 알림 목적지가 삭제됩니다.
stream-logs: 컨테이너 로그를 Dozzle Cloud로 전송
stream-logs-help: AI 기반 분석 및 로그 검색에 필요합니다. 모든 로그 내용을 이 인스턴스에 유지하려면 비활성화하세요.
privacy:
toggle: "로그와 메트릭을 Dozzle Cloud로 전송"
sends-heading: "이 설정이 켜져 있는 동안 컨테이너는 다음을 계속 전송합니다:"
sends-logs: "기록되는 즉시 전송되는 로그 라인"
sends-metrics: "30초 동안 평균한 CPU, 메모리, 네트워크, 디스크 사용량"
never-heading: "전송하지 않는 항목:"
never-env: "환경 변수"
off-note: "로그와 메트릭을 이 인스턴스에만 유지하려면 이 설정을 끄세요. 설정한 알림은 계속 전달되며 원격 컨테이너 접근도 그대로 작동합니다."
per-container: "특정 컨테이너만 제외하려면 라벨을 사용하세요: dev.dozzle.cloud.min_level=disabled"
required-for: "AI 분석, 로그 검색, 리소스 기록에 필요합니다."
welcome:
title: "Cloud는 컨테이너를 위한 지능형 트리아지입니다."
subtitle: "기본 장애 신호 세트가 미리 선택되어 있습니다 — 무엇 때문에 오셨는지 알려 주시면 맞춰 조정할게요."
+10 -2
View File
@@ -359,8 +359,16 @@ cloud:
error-unavailable: Dozzle Cloud is tijdelijk niet beschikbaar. Probeer het later opnieuw.
unlink: Ontkoppelen
unlink-confirm: Weet je zeker dat je de koppeling met Dozzle Cloud wilt opheffen? Dit verwijdert alle cloudmeldingsbestemmingen.
stream-logs: Containerlogs streamen naar Dozzle Cloud
stream-logs-help: Vereist voor AI-gestuurd onderzoek en zoeken in logs. Schakel uit om alle loginhoud op deze instantie te houden.
privacy:
toggle: "Logs en metrieken streamen naar Dozzle Cloud"
sends-heading: "Zolang dit aanstaat, sturen je containers continu:"
sends-logs: "Logregels, zodra ze worden geschreven"
sends-metrics: "CPU-, geheugen-, netwerk- en schijfgebruik, gemiddeld over 30 seconden"
never-heading: "Wordt nooit verstuurd:"
never-env: "Omgevingsvariabelen"
off-note: "Schakel dit uit om logs en metrieken op deze instantie te houden. Ingestelde meldingen worden nog steeds bezorgd en toegang tot containers op afstand blijft werken."
per-container: "Sluit één container uit met het label: dev.dozzle.cloud.min_level=disabled"
required-for: "Vereist voor AI-onderzoek, zoeken in logs en resourcegeschiedenis."
welcome:
title: "Cloud is een slimme triage voor je containers."
subtitle: "Een startset met foutsignalen is voorgeselecteerd — vertel ons wat je hier bracht zodat we het kunnen afstemmen."
+10 -2
View File
@@ -361,8 +361,16 @@ cloud:
error-unavailable: Dozzle Cloud jest chwilowo niedostępny. Spróbuj ponownie później.
unlink: Rozłącz
unlink-confirm: Czy na pewno chcesz rozłączyć się z Dozzle Cloud? Spowoduje to usunięcie wszystkich miejsc docelowych powiadomień w chmurze.
stream-logs: Przesyłaj logi kontenerów do Dozzle Cloud
stream-logs-help: Wymagane do analiz wspieranych przez AI i wyszukiwania w logach. Wyłącz, aby zachować całą zawartość logów w tej instancji.
privacy:
toggle: "Przesyłaj logi i metryki do Dozzle Cloud"
sends-heading: "Gdy ta opcja jest włączona, kontenery stale wysyłają:"
sends-logs: "Wiersze logów, w miarę jak są zapisywane"
sends-metrics: "Zużycie CPU, pamięci, sieci i dysku, uśrednione z 30 sekund"
never-heading: "Nigdy nie jest wysyłane:"
never-env: "Zmienne środowiskowe"
off-note: "Wyłącz tę opcję, aby logi i metryki pozostały w tej instancji. Skonfigurowane alerty nadal będą dostarczane, a zdalny dostęp do kontenerów nadal działa."
per-container: "Wyklucz pojedynczy kontener etykietą: dev.dozzle.cloud.min_level=disabled"
required-for: "Wymagane do analiz AI, wyszukiwania w logach i historii zasobów."
welcome:
title: "Cloud to inteligentna triaż dla Twoich kontenerów."
subtitle: "Startowy zestaw sygnałów awarii jest wstępnie wybrany — powiedz nam, co Cię tu sprowadziło, żebyśmy mogli go dostroić."
+10 -2
View File
@@ -358,8 +358,16 @@ cloud:
error-unavailable: O Dozzle Cloud está temporariamente indisponível. Tente novamente mais tarde.
unlink: Desvincular
unlink-confirm: Tem certeza de que deseja desvincular do Dozzle Cloud? Isso removerá todos os destinos de notificação na nuvem.
stream-logs: Transmitir logs de contêineres para o Dozzle Cloud
stream-logs-help: Necessário para investigações com IA e busca em logs. Desative para manter todo o conteúdo dos logs nesta instância.
privacy:
toggle: "Transmitir logs e métricas para o Dozzle Cloud"
sends-heading: "Enquanto isso estiver ativado, seus contêineres enviam continuamente:"
sends-logs: "Linhas de log, à medida que são escritas"
sends-metrics: "Uso de CPU, memória, rede e disco, com média de 30 segundos"
never-heading: "Nunca é enviado:"
never-env: "Variáveis de ambiente"
off-note: "Desative para manter os logs e as métricas nesta instância. Os alertas que você configurou continuam sendo entregues e o acesso remoto aos contêineres continua funcionando."
per-container: "Exclua um contêiner específico com o rótulo: dev.dozzle.cloud.min_level=disabled"
required-for: "Necessário para investigações com IA, busca em logs e histórico de recursos."
welcome:
title: "Cloud é uma triagem inteligente para os seus contentores."
subtitle: "Um conjunto inicial de sinais de falha está pré-selecionado — conte-nos o que o trouxe aqui para que possamos ajustá-lo."
+10 -2
View File
@@ -358,8 +358,16 @@ cloud:
error-unavailable: Dozzle Cloud временно недоступен. Попробуйте позже.
unlink: Отвязать
unlink-confirm: Вы уверены, что хотите отвязать от Dozzle Cloud? Это удалит все облачные назначения уведомлений.
stream-logs: Передавать логи контейнеров в Dozzle Cloud
stream-logs-help: Требуется для расследований с помощью ИИ и поиска по логам. Отключите, чтобы все содержимое логов оставалось на этом экземпляре.
privacy:
toggle: "Передавать логи и метрики в Dozzle Cloud"
sends-heading: "Пока эта опция включена, ваши контейнеры постоянно отправляют:"
sends-logs: "Строки логов по мере их записи"
sends-metrics: "Использование CPU, памяти, сети и диска, усреднённое за 30 секунд"
never-heading: "Никогда не отправляется:"
never-env: "Переменные окружения"
off-note: "Отключите, чтобы логи и метрики оставались на этом экземпляре. Настроенные оповещения продолжат приходить, а удалённый доступ к контейнерам продолжит работать."
per-container: "Исключить отдельный контейнер меткой: dev.dozzle.cloud.min_level=disabled"
required-for: "Требуется для расследований с помощью ИИ, поиска по логам и истории ресурсов."
welcome:
title: "Cloud — это интеллектуальная триаж-система для ваших контейнеров."
subtitle: "Стартовый набор сигналов сбоев уже выбран — расскажите, что вас привело, чтобы мы могли его настроить."
+10 -2
View File
@@ -364,8 +364,16 @@ cloud:
error-unavailable: Dozzle Cloud je začasno nedosegljiv. Poskusite znova pozneje.
unlink: Prekini povezavo
unlink-confirm: Ali ste prepričani, da želite prekiniti povezavo z Dozzle Cloud? To bo odstranilo vse cilje obvestil v oblaku.
stream-logs: Pretakanje dnevnikov vsebnikov v Dozzle Cloud
stream-logs-help: Potrebno za preiskave z umetno inteligenco in iskanje po dnevnikih. Onemogočite, da vsa vsebina dnevnikov ostane na tej instanci.
privacy:
toggle: "Pretakanje dnevnikov in meritev v Dozzle Cloud"
sends-heading: "Dokler je to vklopljeno, vaši vsebniki nenehno pošiljajo:"
sends-logs: "Vrstice dnevnika, takoj ko so zapisane"
sends-metrics: "Porabo CPE, pomnilnika, omrežja in diska, povprečeno v 30 sekundah"
never-heading: "Nikoli ni poslano:"
never-env: "Spremenljivke okolja"
off-note: "Izklopite, da dnevniki in meritve ostanejo na tej instanci. Nastavljena opozorila se še vedno dostavljajo, oddaljen dostop do vsebnikov pa še naprej deluje."
per-container: "Posamezen vsebnik izključite z oznako: dev.dozzle.cloud.min_level=disabled"
required-for: "Potrebno za preiskave z umetno inteligenco, iskanje po dnevnikih in zgodovino virov."
welcome:
title: "Cloud je inteligentna triaža za vaše vsebnike."
subtitle: "Začetni nabor signalov napak je vnaprej izbran — povejte nam, kaj vas je pripeljalo sem, da ga lahko prilagodimo."
+10 -2
View File
@@ -362,8 +362,16 @@ cloud:
error-unavailable: Dozzle Cloud geçici olarak kullanılamıyor. Lütfen daha sonra tekrar deneyin.
unlink: Bağlantıyı kaldır
unlink-confirm: Dozzle Cloud bağlantısını kaldırmak istediğinizden emin misiniz? Bu, tüm bulut bildirim hedeflerini kaldıracaktır.
stream-logs: Konteyner günlüklerini Dozzle Cloud'a yayınla
stream-logs-help: Yapay zeka destekli incelemeler ve günlük araması için gereklidir. Tüm günlük içeriğini bu örnekte tutmak için devre dışı bırakın.
privacy:
toggle: "Günlükleri ve metrikleri Dozzle Cloud'a yayınla"
sends-heading: "Bu açıkken konteynerleriniz sürekli olarak şunları gönderir:"
sends-logs: "Yazıldıkça günlük satırları"
sends-metrics: "30 saniyelik ortalamayla CPU, bellek, ağ ve disk kullanımı"
never-heading: "Asla gönderilmez:"
never-env: "Ortam değişkenleri"
off-note: "Günlükleri ve metrikleri bu örnekte tutmak için bunu kapatın. Kurduğunuz uyarılar yine iletilir ve uzaktan konteyner erişimi çalışmaya devam eder."
per-container: "Tek bir konteyneri şu etiketle hariç tutun: dev.dozzle.cloud.min_level=disabled"
required-for: "Yapay zeka incelemeleri, günlük araması ve kaynak geçmişi için gereklidir."
welcome:
title: "Cloud, container'larınız için akıllı bir triyajdır."
subtitle: "Başlangıç arıza sinyalleri seti önceden seçili — sizi buraya getiren şeyi söyleyin de ince ayar yapalım."
+10 -2
View File
@@ -361,8 +361,16 @@ cloud:
error-unavailable: Dozzle Cloud 暫時無法使用,請稍後再試。
unlink: 取消連結
unlink-confirm: 確定要取消與 Dozzle Cloud 的連結嗎?這將移除所有雲端通知目標。
stream-logs: 將容器日誌串流至 Dozzle Cloud
stream-logs-help: AI 分析與日誌搜尋所需。停用後,所有日誌內容將保留在此執行個體上。
privacy:
toggle: "將日誌與指標串流至 Dozzle Cloud"
sends-heading: "啟用期間,您的容器會持續傳送:"
sends-logs: "日誌內容,於寫入時傳送"
sends-metrics: "CPU、記憶體、網路與磁碟使用量,以 30 秒為區間取平均"
never-heading: "永不傳送:"
never-env: "環境變數"
off-note: "停用後,日誌與指標將保留在此執行個體上。已設定的警示仍會送達,遠端容器存取也仍可使用。"
per-container: "若要排除單一容器,請使用標籤:dev.dozzle.cloud.min_level=disabled"
required-for: "AI 分析、日誌搜尋與資源歷史記錄所需。"
welcome:
title: "Cloud 是為你的容器提供的智慧分流系統。"
subtitle: "已為你預先選取了一組初始的故障訊號 — 告訴我們你為何而來,我們好做調整。"
+10 -2
View File
@@ -358,8 +358,16 @@ cloud:
error-unavailable: Dozzle Cloud 暂时不可用,请稍后重试。
unlink: 取消关联
unlink-confirm: 确定要取消与 Dozzle Cloud 的关联吗?这将移除所有云通知目标。
stream-logs: 将容器日志流式传输到 Dozzle Cloud
stream-logs-help: AI 调查和日志搜索所需。禁用后,所有日志内容将保留在此实例中。
privacy:
toggle: "将日志和指标流式传输到 Dozzle Cloud"
sends-heading: "启用期间,您的容器会持续发送:"
sends-logs: "日志内容,写入时即发送"
sends-metrics: "CPU、内存、网络和磁盘使用率,按 30 秒取平均"
never-heading: "从不发送:"
never-env: "环境变量"
off-note: "禁用后,日志和指标将保留在此实例中。已配置的告警仍会发送,远程容器访问也仍然可用。"
per-container: "如需排除单个容器,请使用标签:dev.dozzle.cloud.min_level=disabled"
required-for: "AI 调查、日志搜索和资源历史记录所需。"
welcome:
title: "Cloud 是为你的容器提供的智能分诊系统。"
subtitle: "已为你预选了一组初始的故障信号 — 告诉我们你为何而来,我们好做调整。"
+38
View File
@@ -401,6 +401,44 @@ func (l *localCloudHostService) FindContainer(host string, id string, labels con
return nil, fmt.Errorf("host %s not local to this process", host)
}
// SubscribeStats fans stats in from every local client service, stamping the
// originating host on each sample. container.ContainerStat has no host field,
// so without this the cloud side could not tell two same-named containers on
// different hosts apart.
func (l *localCloudHostService) SubscribeStats(ctx context.Context, samples chan<- cloud.StatSample) {
hostIDs := l.resolveHostIDs()
// One inbound channel + forwarder goroutine per service, matching
// SubscribeContainersStarted: a burst on one service must not stall the others.
var dropWarn sync.Once
for i, s := range l.services {
hostID := hostIDs[i]
ch := make(chan container.ContainerStat, 64)
s.SubscribeStats(ctx, ch)
go func() {
for {
select {
case <-ctx.Done():
return
case stat := <-ch:
// Non-blocking on purpose. The stats collector dispatches to
// each subscriber with a blocking send, so if cloud ingest
// ever wedged, backpressure would travel all the way up and
// starve the live UI's stats subscriber on this host. Losing
// a sample only nudges a 30s average; stalling the UI is not
// an acceptable trade for that.
select {
case samples <- cloud.StatSample{Stat: stat, HostID: hostID}:
default:
dropWarn.Do(func() {
log.Warn().Msg("cloud stats: consumer is not keeping up, dropping samples (further drops are silent)")
})
}
}
}
}()
}
}
func (l *localCloudHostService) SubscribeContainersStarted(ctx context.Context, containers chan<- container.Container, filter container_support.ContainerFilter) {
// One inbound channel + forwarder goroutine per service so a slow consumer
// or a burst on one service can't cause the others to drop events.
+380 -122
View File
@@ -196,6 +196,7 @@ type ToolResponse struct {
// *ToolResponse_ListTools
// *ToolResponse_CallTool
// *ToolResponse_LogBatch
// *ToolResponse_StatsBatch
Type isToolResponse_Type `protobuf_oneof:"type"`
unknownFields protoimpl.UnknownFields
sizeCache protoimpl.SizeCache
@@ -272,6 +273,15 @@ func (x *ToolResponse) GetLogBatch() *LogBatch {
return nil
}
func (x *ToolResponse) GetStatsBatch() *StatsBatch {
if x != nil {
if x, ok := x.Type.(*ToolResponse_StatsBatch); ok {
return x.StatsBatch
}
}
return nil
}
type isToolResponse_Type interface {
isToolResponse_Type()
}
@@ -291,12 +301,21 @@ type ToolResponse_LogBatch struct {
LogBatch *LogBatch `protobuf:"bytes,4,opt,name=log_batch,json=logBatch,proto3,oneof"`
}
type ToolResponse_StatsBatch struct {
// Unsolicited server-push: 30s-windowed container resource metrics.
// request_id is empty for stats batches — like log batches, they are not
// replies to a ToolRequest.
StatsBatch *StatsBatch `protobuf:"bytes,5,opt,name=stats_batch,json=statsBatch,proto3,oneof"`
}
func (*ToolResponse_ListTools) isToolResponse_Type() {}
func (*ToolResponse_CallTool) isToolResponse_Type() {}
func (*ToolResponse_LogBatch) isToolResponse_Type() {}
func (*ToolResponse_StatsBatch) isToolResponse_Type() {}
// Batch of log lines from one or more containers. Dozzle pushes these
// continuously while connected. Cloud routes them to VictoriaLogs scoped
// to the owning user (derived from the connection's auth).
@@ -449,6 +468,218 @@ func (x *LogBatchEntry) GetLogId() uint32 {
return 0
}
// Batch of windowed container resource metrics. Dozzle aggregates raw ~1Hz
// stats into fixed windows and pushes one batch per window while connected.
// Cloud writes them to VictoriaMetrics scoped to the owning user (derived
// from the connection's auth).
type StatsBatch struct {
state protoimpl.MessageState `protogen:"open.v1"`
Entries []*StatsBatchEntry `protobuf:"bytes,1,rep,name=entries,proto3" json:"entries,omitempty"`
unknownFields protoimpl.UnknownFields
sizeCache protoimpl.SizeCache
}
func (x *StatsBatch) Reset() {
*x = StatsBatch{}
mi := &file_cloud_proto_msgTypes[4]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
func (x *StatsBatch) String() string {
return protoimpl.X.MessageStringOf(x)
}
func (*StatsBatch) ProtoMessage() {}
func (x *StatsBatch) ProtoReflect() protoreflect.Message {
mi := &file_cloud_proto_msgTypes[4]
if x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
ms.StoreMessageInfo(mi)
}
return ms
}
return mi.MessageOf(x)
}
// Deprecated: Use StatsBatch.ProtoReflect.Descriptor instead.
func (*StatsBatch) Descriptor() ([]byte, []int) {
return file_cloud_proto_rawDescGZIP(), []int{4}
}
func (x *StatsBatch) GetEntries() []*StatsBatchEntry {
if x != nil {
return x.Entries
}
return nil
}
// One container's metrics for one aggregation window.
//
// Series identity on the Cloud side is (host_id, container_name) — there is
// deliberately no container_id. Container ids churn on every redeploy and each
// new value would strand the old time series forever; keying on the name means
// a redeploy continues the same series, which is what a resource chart wants.
type StatsBatchEntry struct {
state protoimpl.MessageState `protogen:"open.v1"`
HostId string `protobuf:"bytes,1,opt,name=host_id,json=hostId,proto3" json:"host_id,omitempty"`
ContainerName string `protobuf:"bytes,2,opt,name=container_name,json=containerName,proto3" json:"container_name,omitempty"`
TimestampNs int64 `protobuf:"varint,3,opt,name=timestamp_ns,json=timestampNs,proto3" json:"timestamp_ns,omitempty"` // window END, unix nanoseconds
// Gauges — arithmetic mean over the window.
CpuPercent float64 `protobuf:"fixed64,4,opt,name=cpu_percent,json=cpuPercent,proto3" json:"cpu_percent,omitempty"` // normalised 0-100 across cores, NOT per-core
MemoryPercent float64 `protobuf:"fixed64,5,opt,name=memory_percent,json=memoryPercent,proto3" json:"memory_percent,omitempty"`
MemoryUsageBytes float64 `protobuf:"fixed64,6,opt,name=memory_usage_bytes,json=memoryUsageBytes,proto3" json:"memory_usage_bytes,omitempty"`
// Peak sample in the window. A 30s mean flattens spikes, and spikes are
// exactly what an incident chart is read for.
CpuPercentMax float64 `protobuf:"fixed64,7,opt,name=cpu_percent_max,json=cpuPercentMax,proto3" json:"cpu_percent_max,omitempty"`
MemoryPercentMax float64 `protobuf:"fixed64,8,opt,name=memory_percent_max,json=memoryPercentMax,proto3" json:"memory_percent_max,omitempty"`
// Monotonic counters — the LAST value observed in the window, not a delta.
// Cloud stores these as Prometheus counters so PromQL rate()/increase()
// handle the reset that happens when a container is recreated.
NetworkRxTotal uint64 `protobuf:"varint,9,opt,name=network_rx_total,json=networkRxTotal,proto3" json:"network_rx_total,omitempty"`
NetworkTxTotal uint64 `protobuf:"varint,10,opt,name=network_tx_total,json=networkTxTotal,proto3" json:"network_tx_total,omitempty"`
DiskReadTotal uint64 `protobuf:"varint,11,opt,name=disk_read_total,json=diskReadTotal,proto3" json:"disk_read_total,omitempty"`
DiskWriteTotal uint64 `protobuf:"varint,12,opt,name=disk_write_total,json=diskWriteTotal,proto3" json:"disk_write_total,omitempty"`
// Number of raw samples folded into this entry. Windows with zero samples
// are never emitted.
Samples uint32 `protobuf:"varint,13,opt,name=samples,proto3" json:"samples,omitempty"`
// The core count cpu_percent was divided by, so the raw per-core figure
// Docker reports stays recoverable downstream.
CpuCores float64 `protobuf:"fixed64,14,opt,name=cpu_cores,json=cpuCores,proto3" json:"cpu_cores,omitempty"`
unknownFields protoimpl.UnknownFields
sizeCache protoimpl.SizeCache
}
func (x *StatsBatchEntry) Reset() {
*x = StatsBatchEntry{}
mi := &file_cloud_proto_msgTypes[5]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
func (x *StatsBatchEntry) String() string {
return protoimpl.X.MessageStringOf(x)
}
func (*StatsBatchEntry) ProtoMessage() {}
func (x *StatsBatchEntry) ProtoReflect() protoreflect.Message {
mi := &file_cloud_proto_msgTypes[5]
if x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
ms.StoreMessageInfo(mi)
}
return ms
}
return mi.MessageOf(x)
}
// Deprecated: Use StatsBatchEntry.ProtoReflect.Descriptor instead.
func (*StatsBatchEntry) Descriptor() ([]byte, []int) {
return file_cloud_proto_rawDescGZIP(), []int{5}
}
func (x *StatsBatchEntry) GetHostId() string {
if x != nil {
return x.HostId
}
return ""
}
func (x *StatsBatchEntry) GetContainerName() string {
if x != nil {
return x.ContainerName
}
return ""
}
func (x *StatsBatchEntry) GetTimestampNs() int64 {
if x != nil {
return x.TimestampNs
}
return 0
}
func (x *StatsBatchEntry) GetCpuPercent() float64 {
if x != nil {
return x.CpuPercent
}
return 0
}
func (x *StatsBatchEntry) GetMemoryPercent() float64 {
if x != nil {
return x.MemoryPercent
}
return 0
}
func (x *StatsBatchEntry) GetMemoryUsageBytes() float64 {
if x != nil {
return x.MemoryUsageBytes
}
return 0
}
func (x *StatsBatchEntry) GetCpuPercentMax() float64 {
if x != nil {
return x.CpuPercentMax
}
return 0
}
func (x *StatsBatchEntry) GetMemoryPercentMax() float64 {
if x != nil {
return x.MemoryPercentMax
}
return 0
}
func (x *StatsBatchEntry) GetNetworkRxTotal() uint64 {
if x != nil {
return x.NetworkRxTotal
}
return 0
}
func (x *StatsBatchEntry) GetNetworkTxTotal() uint64 {
if x != nil {
return x.NetworkTxTotal
}
return 0
}
func (x *StatsBatchEntry) GetDiskReadTotal() uint64 {
if x != nil {
return x.DiskReadTotal
}
return 0
}
func (x *StatsBatchEntry) GetDiskWriteTotal() uint64 {
if x != nil {
return x.DiskWriteTotal
}
return 0
}
func (x *StatsBatchEntry) GetSamples() uint32 {
if x != nil {
return x.Samples
}
return 0
}
func (x *StatsBatchEntry) GetCpuCores() float64 {
if x != nil {
return x.CpuCores
}
return 0
}
type ListToolsRequest struct {
state protoimpl.MessageState `protogen:"open.v1"`
unknownFields protoimpl.UnknownFields
@@ -457,7 +688,7 @@ type ListToolsRequest struct {
func (x *ListToolsRequest) Reset() {
*x = ListToolsRequest{}
mi := &file_cloud_proto_msgTypes[4]
mi := &file_cloud_proto_msgTypes[6]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
@@ -469,7 +700,7 @@ func (x *ListToolsRequest) String() string {
func (*ListToolsRequest) ProtoMessage() {}
func (x *ListToolsRequest) ProtoReflect() protoreflect.Message {
mi := &file_cloud_proto_msgTypes[4]
mi := &file_cloud_proto_msgTypes[6]
if x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
@@ -482,7 +713,7 @@ func (x *ListToolsRequest) ProtoReflect() protoreflect.Message {
// Deprecated: Use ListToolsRequest.ProtoReflect.Descriptor instead.
func (*ListToolsRequest) Descriptor() ([]byte, []int) {
return file_cloud_proto_rawDescGZIP(), []int{4}
return file_cloud_proto_rawDescGZIP(), []int{6}
}
type ListToolsResponse struct {
@@ -495,7 +726,7 @@ type ListToolsResponse struct {
func (x *ListToolsResponse) Reset() {
*x = ListToolsResponse{}
mi := &file_cloud_proto_msgTypes[5]
mi := &file_cloud_proto_msgTypes[7]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
@@ -507,7 +738,7 @@ func (x *ListToolsResponse) String() string {
func (*ListToolsResponse) ProtoMessage() {}
func (x *ListToolsResponse) ProtoReflect() protoreflect.Message {
mi := &file_cloud_proto_msgTypes[5]
mi := &file_cloud_proto_msgTypes[7]
if x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
@@ -520,7 +751,7 @@ func (x *ListToolsResponse) ProtoReflect() protoreflect.Message {
// Deprecated: Use ListToolsResponse.ProtoReflect.Descriptor instead.
func (*ListToolsResponse) Descriptor() ([]byte, []int) {
return file_cloud_proto_rawDescGZIP(), []int{5}
return file_cloud_proto_rawDescGZIP(), []int{7}
}
func (x *ListToolsResponse) GetTools() []*ToolDefinition {
@@ -562,7 +793,7 @@ type ToolDefinition struct {
func (x *ToolDefinition) Reset() {
*x = ToolDefinition{}
mi := &file_cloud_proto_msgTypes[6]
mi := &file_cloud_proto_msgTypes[8]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
@@ -574,7 +805,7 @@ func (x *ToolDefinition) String() string {
func (*ToolDefinition) ProtoMessage() {}
func (x *ToolDefinition) ProtoReflect() protoreflect.Message {
mi := &file_cloud_proto_msgTypes[6]
mi := &file_cloud_proto_msgTypes[8]
if x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
@@ -587,7 +818,7 @@ func (x *ToolDefinition) ProtoReflect() protoreflect.Message {
// Deprecated: Use ToolDefinition.ProtoReflect.Descriptor instead.
func (*ToolDefinition) Descriptor() ([]byte, []int) {
return file_cloud_proto_rawDescGZIP(), []int{6}
return file_cloud_proto_rawDescGZIP(), []int{8}
}
func (x *ToolDefinition) GetName() string {
@@ -635,7 +866,7 @@ type CallToolRequest struct {
func (x *CallToolRequest) Reset() {
*x = CallToolRequest{}
mi := &file_cloud_proto_msgTypes[7]
mi := &file_cloud_proto_msgTypes[9]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
@@ -647,7 +878,7 @@ func (x *CallToolRequest) String() string {
func (*CallToolRequest) ProtoMessage() {}
func (x *CallToolRequest) ProtoReflect() protoreflect.Message {
mi := &file_cloud_proto_msgTypes[7]
mi := &file_cloud_proto_msgTypes[9]
if x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
@@ -660,7 +891,7 @@ func (x *CallToolRequest) ProtoReflect() protoreflect.Message {
// Deprecated: Use CallToolRequest.ProtoReflect.Descriptor instead.
func (*CallToolRequest) Descriptor() ([]byte, []int) {
return file_cloud_proto_rawDescGZIP(), []int{7}
return file_cloud_proto_rawDescGZIP(), []int{9}
}
func (x *CallToolRequest) GetName() string {
@@ -700,7 +931,7 @@ type CallToolResponse struct {
func (x *CallToolResponse) Reset() {
*x = CallToolResponse{}
mi := &file_cloud_proto_msgTypes[8]
mi := &file_cloud_proto_msgTypes[10]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
@@ -712,7 +943,7 @@ func (x *CallToolResponse) String() string {
func (*CallToolResponse) ProtoMessage() {}
func (x *CallToolResponse) ProtoReflect() protoreflect.Message {
mi := &file_cloud_proto_msgTypes[8]
mi := &file_cloud_proto_msgTypes[10]
if x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
@@ -725,7 +956,7 @@ func (x *CallToolResponse) ProtoReflect() protoreflect.Message {
// Deprecated: Use CallToolResponse.ProtoReflect.Descriptor instead.
func (*CallToolResponse) Descriptor() ([]byte, []int) {
return file_cloud_proto_rawDescGZIP(), []int{8}
return file_cloud_proto_rawDescGZIP(), []int{10}
}
func (x *CallToolResponse) GetSuccess() bool {
@@ -896,7 +1127,7 @@ type CancelStreamRequest struct {
func (x *CancelStreamRequest) Reset() {
*x = CancelStreamRequest{}
mi := &file_cloud_proto_msgTypes[9]
mi := &file_cloud_proto_msgTypes[11]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
@@ -908,7 +1139,7 @@ func (x *CancelStreamRequest) String() string {
func (*CancelStreamRequest) ProtoMessage() {}
func (x *CancelStreamRequest) ProtoReflect() protoreflect.Message {
mi := &file_cloud_proto_msgTypes[9]
mi := &file_cloud_proto_msgTypes[11]
if x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
@@ -921,7 +1152,7 @@ func (x *CancelStreamRequest) ProtoReflect() protoreflect.Message {
// Deprecated: Use CancelStreamRequest.ProtoReflect.Descriptor instead.
func (*CancelStreamRequest) Descriptor() ([]byte, []int) {
return file_cloud_proto_rawDescGZIP(), []int{9}
return file_cloud_proto_rawDescGZIP(), []int{11}
}
func (x *CancelStreamRequest) GetStreamRequestId() string {
@@ -948,7 +1179,7 @@ type HostInfo struct {
func (x *HostInfo) Reset() {
*x = HostInfo{}
mi := &file_cloud_proto_msgTypes[10]
mi := &file_cloud_proto_msgTypes[12]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
@@ -960,7 +1191,7 @@ func (x *HostInfo) String() string {
func (*HostInfo) ProtoMessage() {}
func (x *HostInfo) ProtoReflect() protoreflect.Message {
mi := &file_cloud_proto_msgTypes[10]
mi := &file_cloud_proto_msgTypes[12]
if x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
@@ -973,7 +1204,7 @@ func (x *HostInfo) ProtoReflect() protoreflect.Message {
// Deprecated: Use HostInfo.ProtoReflect.Descriptor instead.
func (*HostInfo) Descriptor() ([]byte, []int) {
return file_cloud_proto_rawDescGZIP(), []int{10}
return file_cloud_proto_rawDescGZIP(), []int{12}
}
func (x *HostInfo) GetId() string {
@@ -1041,7 +1272,7 @@ type ListHostsResult struct {
func (x *ListHostsResult) Reset() {
*x = ListHostsResult{}
mi := &file_cloud_proto_msgTypes[11]
mi := &file_cloud_proto_msgTypes[13]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
@@ -1053,7 +1284,7 @@ func (x *ListHostsResult) String() string {
func (*ListHostsResult) ProtoMessage() {}
func (x *ListHostsResult) ProtoReflect() protoreflect.Message {
mi := &file_cloud_proto_msgTypes[11]
mi := &file_cloud_proto_msgTypes[13]
if x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
@@ -1066,7 +1297,7 @@ func (x *ListHostsResult) ProtoReflect() protoreflect.Message {
// Deprecated: Use ListHostsResult.ProtoReflect.Descriptor instead.
func (*ListHostsResult) Descriptor() ([]byte, []int) {
return file_cloud_proto_rawDescGZIP(), []int{11}
return file_cloud_proto_rawDescGZIP(), []int{13}
}
func (x *ListHostsResult) GetHosts() []*HostInfo {
@@ -1097,7 +1328,7 @@ type ContainerInfo struct {
func (x *ContainerInfo) Reset() {
*x = ContainerInfo{}
mi := &file_cloud_proto_msgTypes[12]
mi := &file_cloud_proto_msgTypes[14]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
@@ -1109,7 +1340,7 @@ func (x *ContainerInfo) String() string {
func (*ContainerInfo) ProtoMessage() {}
func (x *ContainerInfo) ProtoReflect() protoreflect.Message {
mi := &file_cloud_proto_msgTypes[12]
mi := &file_cloud_proto_msgTypes[14]
if x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
@@ -1122,7 +1353,7 @@ func (x *ContainerInfo) ProtoReflect() protoreflect.Message {
// Deprecated: Use ContainerInfo.ProtoReflect.Descriptor instead.
func (*ContainerInfo) Descriptor() ([]byte, []int) {
return file_cloud_proto_rawDescGZIP(), []int{12}
return file_cloud_proto_rawDescGZIP(), []int{14}
}
func (x *ContainerInfo) GetId() string {
@@ -1218,7 +1449,7 @@ type ListContainersResult struct {
func (x *ListContainersResult) Reset() {
*x = ListContainersResult{}
mi := &file_cloud_proto_msgTypes[13]
mi := &file_cloud_proto_msgTypes[15]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
@@ -1230,7 +1461,7 @@ func (x *ListContainersResult) String() string {
func (*ListContainersResult) ProtoMessage() {}
func (x *ListContainersResult) ProtoReflect() protoreflect.Message {
mi := &file_cloud_proto_msgTypes[13]
mi := &file_cloud_proto_msgTypes[15]
if x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
@@ -1243,7 +1474,7 @@ func (x *ListContainersResult) ProtoReflect() protoreflect.Message {
// Deprecated: Use ListContainersResult.ProtoReflect.Descriptor instead.
func (*ListContainersResult) Descriptor() ([]byte, []int) {
return file_cloud_proto_rawDescGZIP(), []int{13}
return file_cloud_proto_rawDescGZIP(), []int{15}
}
func (x *ListContainersResult) GetContainers() []*ContainerInfo {
@@ -1275,7 +1506,7 @@ type ContainerStatEntry struct {
func (x *ContainerStatEntry) Reset() {
*x = ContainerStatEntry{}
mi := &file_cloud_proto_msgTypes[14]
mi := &file_cloud_proto_msgTypes[16]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
@@ -1287,7 +1518,7 @@ func (x *ContainerStatEntry) String() string {
func (*ContainerStatEntry) ProtoMessage() {}
func (x *ContainerStatEntry) ProtoReflect() protoreflect.Message {
mi := &file_cloud_proto_msgTypes[14]
mi := &file_cloud_proto_msgTypes[16]
if x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
@@ -1300,7 +1531,7 @@ func (x *ContainerStatEntry) ProtoReflect() protoreflect.Message {
// Deprecated: Use ContainerStatEntry.ProtoReflect.Descriptor instead.
func (*ContainerStatEntry) Descriptor() ([]byte, []int) {
return file_cloud_proto_rawDescGZIP(), []int{14}
return file_cloud_proto_rawDescGZIP(), []int{16}
}
func (x *ContainerStatEntry) GetId() string {
@@ -1403,7 +1634,7 @@ type ContainerStatsResult struct {
func (x *ContainerStatsResult) Reset() {
*x = ContainerStatsResult{}
mi := &file_cloud_proto_msgTypes[15]
mi := &file_cloud_proto_msgTypes[17]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
@@ -1415,7 +1646,7 @@ func (x *ContainerStatsResult) String() string {
func (*ContainerStatsResult) ProtoMessage() {}
func (x *ContainerStatsResult) ProtoReflect() protoreflect.Message {
mi := &file_cloud_proto_msgTypes[15]
mi := &file_cloud_proto_msgTypes[17]
if x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
@@ -1428,7 +1659,7 @@ func (x *ContainerStatsResult) ProtoReflect() protoreflect.Message {
// Deprecated: Use ContainerStatsResult.ProtoReflect.Descriptor instead.
func (*ContainerStatsResult) Descriptor() ([]byte, []int) {
return file_cloud_proto_rawDescGZIP(), []int{15}
return file_cloud_proto_rawDescGZIP(), []int{17}
}
func (x *ContainerStatsResult) GetStats() []*ContainerStatEntry {
@@ -1451,7 +1682,7 @@ type LogEntry struct {
func (x *LogEntry) Reset() {
*x = LogEntry{}
mi := &file_cloud_proto_msgTypes[16]
mi := &file_cloud_proto_msgTypes[18]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
@@ -1463,7 +1694,7 @@ func (x *LogEntry) String() string {
func (*LogEntry) ProtoMessage() {}
func (x *LogEntry) ProtoReflect() protoreflect.Message {
mi := &file_cloud_proto_msgTypes[16]
mi := &file_cloud_proto_msgTypes[18]
if x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
@@ -1476,7 +1707,7 @@ func (x *LogEntry) ProtoReflect() protoreflect.Message {
// Deprecated: Use LogEntry.ProtoReflect.Descriptor instead.
func (*LogEntry) Descriptor() ([]byte, []int) {
return file_cloud_proto_rawDescGZIP(), []int{16}
return file_cloud_proto_rawDescGZIP(), []int{18}
}
func (x *LogEntry) GetTimestamp() int64 {
@@ -1517,7 +1748,7 @@ type FetchLogsResult struct {
func (x *FetchLogsResult) Reset() {
*x = FetchLogsResult{}
mi := &file_cloud_proto_msgTypes[17]
mi := &file_cloud_proto_msgTypes[19]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
@@ -1529,7 +1760,7 @@ func (x *FetchLogsResult) String() string {
func (*FetchLogsResult) ProtoMessage() {}
func (x *FetchLogsResult) ProtoReflect() protoreflect.Message {
mi := &file_cloud_proto_msgTypes[17]
mi := &file_cloud_proto_msgTypes[19]
if x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
@@ -1542,7 +1773,7 @@ func (x *FetchLogsResult) ProtoReflect() protoreflect.Message {
// Deprecated: Use FetchLogsResult.ProtoReflect.Descriptor instead.
func (*FetchLogsResult) Descriptor() ([]byte, []int) {
return file_cloud_proto_rawDescGZIP(), []int{17}
return file_cloud_proto_rawDescGZIP(), []int{19}
}
func (x *FetchLogsResult) GetContainerName() string {
@@ -1586,7 +1817,7 @@ type InspectContainerResult struct {
func (x *InspectContainerResult) Reset() {
*x = InspectContainerResult{}
mi := &file_cloud_proto_msgTypes[18]
mi := &file_cloud_proto_msgTypes[20]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
@@ -1598,7 +1829,7 @@ func (x *InspectContainerResult) String() string {
func (*InspectContainerResult) ProtoMessage() {}
func (x *InspectContainerResult) ProtoReflect() protoreflect.Message {
mi := &file_cloud_proto_msgTypes[18]
mi := &file_cloud_proto_msgTypes[20]
if x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
@@ -1611,7 +1842,7 @@ func (x *InspectContainerResult) ProtoReflect() protoreflect.Message {
// Deprecated: Use InspectContainerResult.ProtoReflect.Descriptor instead.
func (*InspectContainerResult) Descriptor() ([]byte, []int) {
return file_cloud_proto_rawDescGZIP(), []int{18}
return file_cloud_proto_rawDescGZIP(), []int{20}
}
func (x *InspectContainerResult) GetId() string {
@@ -1753,7 +1984,7 @@ type ActionResult struct {
func (x *ActionResult) Reset() {
*x = ActionResult{}
mi := &file_cloud_proto_msgTypes[19]
mi := &file_cloud_proto_msgTypes[21]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
@@ -1765,7 +1996,7 @@ func (x *ActionResult) String() string {
func (*ActionResult) ProtoMessage() {}
func (x *ActionResult) ProtoReflect() protoreflect.Message {
mi := &file_cloud_proto_msgTypes[19]
mi := &file_cloud_proto_msgTypes[21]
if x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
@@ -1778,7 +2009,7 @@ func (x *ActionResult) ProtoReflect() protoreflect.Message {
// Deprecated: Use ActionResult.ProtoReflect.Descriptor instead.
func (*ActionResult) Descriptor() ([]byte, []int) {
return file_cloud_proto_rawDescGZIP(), []int{19}
return file_cloud_proto_rawDescGZIP(), []int{21}
}
func (x *ActionResult) GetSuccess() bool {
@@ -1821,7 +2052,7 @@ type DeployResult struct {
func (x *DeployResult) Reset() {
*x = DeployResult{}
mi := &file_cloud_proto_msgTypes[20]
mi := &file_cloud_proto_msgTypes[22]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
@@ -1833,7 +2064,7 @@ func (x *DeployResult) String() string {
func (*DeployResult) ProtoMessage() {}
func (x *DeployResult) ProtoReflect() protoreflect.Message {
mi := &file_cloud_proto_msgTypes[20]
mi := &file_cloud_proto_msgTypes[22]
if x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
@@ -1846,7 +2077,7 @@ func (x *DeployResult) ProtoReflect() protoreflect.Message {
// Deprecated: Use DeployResult.ProtoReflect.Descriptor instead.
func (*DeployResult) Descriptor() ([]byte, []int) {
return file_cloud_proto_rawDescGZIP(), []int{20}
return file_cloud_proto_rawDescGZIP(), []int{22}
}
func (x *DeployResult) GetSuccess() bool {
@@ -1881,7 +2112,7 @@ type NotificationResult struct {
func (x *NotificationResult) Reset() {
*x = NotificationResult{}
mi := &file_cloud_proto_msgTypes[21]
mi := &file_cloud_proto_msgTypes[23]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
@@ -1893,7 +2124,7 @@ func (x *NotificationResult) String() string {
func (*NotificationResult) ProtoMessage() {}
func (x *NotificationResult) ProtoReflect() protoreflect.Message {
mi := &file_cloud_proto_msgTypes[21]
mi := &file_cloud_proto_msgTypes[23]
if x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
@@ -1906,7 +2137,7 @@ func (x *NotificationResult) ProtoReflect() protoreflect.Message {
// Deprecated: Use NotificationResult.ProtoReflect.Descriptor instead.
func (*NotificationResult) Descriptor() ([]byte, []int) {
return file_cloud_proto_rawDescGZIP(), []int{21}
return file_cloud_proto_rawDescGZIP(), []int{23}
}
func (x *NotificationResult) GetSuccess() bool {
@@ -1945,7 +2176,7 @@ type SearchLogsRequest struct {
func (x *SearchLogsRequest) Reset() {
*x = SearchLogsRequest{}
mi := &file_cloud_proto_msgTypes[22]
mi := &file_cloud_proto_msgTypes[24]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
@@ -1957,7 +2188,7 @@ func (x *SearchLogsRequest) String() string {
func (*SearchLogsRequest) ProtoMessage() {}
func (x *SearchLogsRequest) ProtoReflect() protoreflect.Message {
mi := &file_cloud_proto_msgTypes[22]
mi := &file_cloud_proto_msgTypes[24]
if x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
@@ -1970,7 +2201,7 @@ func (x *SearchLogsRequest) ProtoReflect() protoreflect.Message {
// Deprecated: Use SearchLogsRequest.ProtoReflect.Descriptor instead.
func (*SearchLogsRequest) Descriptor() ([]byte, []int) {
return file_cloud_proto_rawDescGZIP(), []int{22}
return file_cloud_proto_rawDescGZIP(), []int{24}
}
func (x *SearchLogsRequest) GetQuery() string {
@@ -2020,7 +2251,7 @@ type SearchLogsResponse struct {
func (x *SearchLogsResponse) Reset() {
*x = SearchLogsResponse{}
mi := &file_cloud_proto_msgTypes[23]
mi := &file_cloud_proto_msgTypes[25]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
@@ -2032,7 +2263,7 @@ func (x *SearchLogsResponse) String() string {
func (*SearchLogsResponse) ProtoMessage() {}
func (x *SearchLogsResponse) ProtoReflect() protoreflect.Message {
mi := &file_cloud_proto_msgTypes[23]
mi := &file_cloud_proto_msgTypes[25]
if x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
@@ -2045,7 +2276,7 @@ func (x *SearchLogsResponse) ProtoReflect() protoreflect.Message {
// Deprecated: Use SearchLogsResponse.ProtoReflect.Descriptor instead.
func (*SearchLogsResponse) Descriptor() ([]byte, []int) {
return file_cloud_proto_rawDescGZIP(), []int{23}
return file_cloud_proto_rawDescGZIP(), []int{25}
}
func (x *SearchLogsResponse) GetHits() []*SearchLogHit {
@@ -2089,7 +2320,7 @@ type SearchLogHit struct {
func (x *SearchLogHit) Reset() {
*x = SearchLogHit{}
mi := &file_cloud_proto_msgTypes[24]
mi := &file_cloud_proto_msgTypes[26]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
@@ -2101,7 +2332,7 @@ func (x *SearchLogHit) String() string {
func (*SearchLogHit) ProtoMessage() {}
func (x *SearchLogHit) ProtoReflect() protoreflect.Message {
mi := &file_cloud_proto_msgTypes[24]
mi := &file_cloud_proto_msgTypes[26]
if x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
@@ -2114,7 +2345,7 @@ func (x *SearchLogHit) ProtoReflect() protoreflect.Message {
// Deprecated: Use SearchLogHit.ProtoReflect.Descriptor instead.
func (*SearchLogHit) Descriptor() ([]byte, []int) {
return file_cloud_proto_rawDescGZIP(), []int{24}
return file_cloud_proto_rawDescGZIP(), []int{26}
}
func (x *SearchLogHit) GetTimestampNs() int64 {
@@ -2185,14 +2416,16 @@ const file_cloud_proto_rawDesc = "" +
"list_tools\x18\x02 \x01(\v2\x17.cloud.ListToolsRequestH\x00R\tlistTools\x125\n" +
"\tcall_tool\x18\x03 \x01(\v2\x16.cloud.CallToolRequestH\x00R\bcallTool\x12A\n" +
"\rcancel_stream\x18\x04 \x01(\v2\x1a.cloud.CancelStreamRequestH\x00R\fcancelStreamB\x06\n" +
"\x04type\"\xd8\x01\n" +
"\x04type\"\x8e\x02\n" +
"\fToolResponse\x12\x1d\n" +
"\n" +
"request_id\x18\x01 \x01(\tR\trequestId\x129\n" +
"\n" +
"list_tools\x18\x02 \x01(\v2\x18.cloud.ListToolsResponseH\x00R\tlistTools\x126\n" +
"\tcall_tool\x18\x03 \x01(\v2\x17.cloud.CallToolResponseH\x00R\bcallTool\x12.\n" +
"\tlog_batch\x18\x04 \x01(\v2\x0f.cloud.LogBatchH\x00R\blogBatchB\x06\n" +
"\tlog_batch\x18\x04 \x01(\v2\x0f.cloud.LogBatchH\x00R\blogBatch\x124\n" +
"\vstats_batch\x18\x05 \x01(\v2\x11.cloud.StatsBatchH\x00R\n" +
"statsBatchB\x06\n" +
"\x04type\":\n" +
"\bLogBatch\x12.\n" +
"\aentries\x18\x01 \x03(\v2\x14.cloud.LogBatchEntryR\aentries\"\xf4\x01\n" +
@@ -2204,7 +2437,27 @@ const file_cloud_proto_rawDesc = "" +
"\amessage\x18\x05 \x01(\tR\amessage\x12\x16\n" +
"\x06stream\x18\x06 \x01(\tR\x06stream\x12\x14\n" +
"\x05level\x18\a \x01(\tR\x05level\x12\x15\n" +
"\x06log_id\x18\b \x01(\rR\x05logId\"\x12\n" +
"\x06log_id\x18\b \x01(\rR\x05logId\">\n" +
"\n" +
"StatsBatch\x120\n" +
"\aentries\x18\x01 \x03(\v2\x16.cloud.StatsBatchEntryR\aentries\"\x9d\x04\n" +
"\x0fStatsBatchEntry\x12\x17\n" +
"\ahost_id\x18\x01 \x01(\tR\x06hostId\x12%\n" +
"\x0econtainer_name\x18\x02 \x01(\tR\rcontainerName\x12!\n" +
"\ftimestamp_ns\x18\x03 \x01(\x03R\vtimestampNs\x12\x1f\n" +
"\vcpu_percent\x18\x04 \x01(\x01R\n" +
"cpuPercent\x12%\n" +
"\x0ememory_percent\x18\x05 \x01(\x01R\rmemoryPercent\x12,\n" +
"\x12memory_usage_bytes\x18\x06 \x01(\x01R\x10memoryUsageBytes\x12&\n" +
"\x0fcpu_percent_max\x18\a \x01(\x01R\rcpuPercentMax\x12,\n" +
"\x12memory_percent_max\x18\b \x01(\x01R\x10memoryPercentMax\x12(\n" +
"\x10network_rx_total\x18\t \x01(\x04R\x0enetworkRxTotal\x12(\n" +
"\x10network_tx_total\x18\n" +
" \x01(\x04R\x0enetworkTxTotal\x12&\n" +
"\x0fdisk_read_total\x18\v \x01(\x04R\rdiskReadTotal\x12(\n" +
"\x10disk_write_total\x18\f \x01(\x04R\x0ediskWriteTotal\x12\x18\n" +
"\asamples\x18\r \x01(\rR\asamples\x12\x1b\n" +
"\tcpu_cores\x18\x0e \x01(\x01R\bcpuCores\"\x12\n" +
"\x10ListToolsRequest\"Z\n" +
"\x11ListToolsResponse\x12+\n" +
"\x05tools\x18\x01 \x03(\v2\x15.cloud.ToolDefinitionR\x05tools\x12\x18\n" +
@@ -2377,69 +2630,73 @@ func file_cloud_proto_rawDescGZIP() []byte {
}
var file_cloud_proto_enumTypes = make([]protoimpl.EnumInfo, 1)
var file_cloud_proto_msgTypes = make([]protoimpl.MessageInfo, 26)
var file_cloud_proto_msgTypes = make([]protoimpl.MessageInfo, 28)
var file_cloud_proto_goTypes = []any{
(ToolScope)(0), // 0: cloud.ToolScope
(*ToolRequest)(nil), // 1: cloud.ToolRequest
(*ToolResponse)(nil), // 2: cloud.ToolResponse
(*LogBatch)(nil), // 3: cloud.LogBatch
(*LogBatchEntry)(nil), // 4: cloud.LogBatchEntry
(*ListToolsRequest)(nil), // 5: cloud.ListToolsRequest
(*ListToolsResponse)(nil), // 6: cloud.ListToolsResponse
(*ToolDefinition)(nil), // 7: cloud.ToolDefinition
(*CallToolRequest)(nil), // 8: cloud.CallToolRequest
(*CallToolResponse)(nil), // 9: cloud.CallToolResponse
(*CancelStreamRequest)(nil), // 10: cloud.CancelStreamRequest
(*HostInfo)(nil), // 11: cloud.HostInfo
(*ListHostsResult)(nil), // 12: cloud.ListHostsResult
(*ContainerInfo)(nil), // 13: cloud.ContainerInfo
(*ListContainersResult)(nil), // 14: cloud.ListContainersResult
(*ContainerStatEntry)(nil), // 15: cloud.ContainerStatEntry
(*ContainerStatsResult)(nil), // 16: cloud.ContainerStatsResult
(*LogEntry)(nil), // 17: cloud.LogEntry
(*FetchLogsResult)(nil), // 18: cloud.FetchLogsResult
(*InspectContainerResult)(nil), // 19: cloud.InspectContainerResult
(*ActionResult)(nil), // 20: cloud.ActionResult
(*DeployResult)(nil), // 21: cloud.DeployResult
(*NotificationResult)(nil), // 22: cloud.NotificationResult
(*SearchLogsRequest)(nil), // 23: cloud.SearchLogsRequest
(*SearchLogsResponse)(nil), // 24: cloud.SearchLogsResponse
(*SearchLogHit)(nil), // 25: cloud.SearchLogHit
nil, // 26: cloud.InspectContainerResult.LabelsEntry
(*StatsBatch)(nil), // 5: cloud.StatsBatch
(*StatsBatchEntry)(nil), // 6: cloud.StatsBatchEntry
(*ListToolsRequest)(nil), // 7: cloud.ListToolsRequest
(*ListToolsResponse)(nil), // 8: cloud.ListToolsResponse
(*ToolDefinition)(nil), // 9: cloud.ToolDefinition
(*CallToolRequest)(nil), // 10: cloud.CallToolRequest
(*CallToolResponse)(nil), // 11: cloud.CallToolResponse
(*CancelStreamRequest)(nil), // 12: cloud.CancelStreamRequest
(*HostInfo)(nil), // 13: cloud.HostInfo
(*ListHostsResult)(nil), // 14: cloud.ListHostsResult
(*ContainerInfo)(nil), // 15: cloud.ContainerInfo
(*ListContainersResult)(nil), // 16: cloud.ListContainersResult
(*ContainerStatEntry)(nil), // 17: cloud.ContainerStatEntry
(*ContainerStatsResult)(nil), // 18: cloud.ContainerStatsResult
(*LogEntry)(nil), // 19: cloud.LogEntry
(*FetchLogsResult)(nil), // 20: cloud.FetchLogsResult
(*InspectContainerResult)(nil), // 21: cloud.InspectContainerResult
(*ActionResult)(nil), // 22: cloud.ActionResult
(*DeployResult)(nil), // 23: cloud.DeployResult
(*NotificationResult)(nil), // 24: cloud.NotificationResult
(*SearchLogsRequest)(nil), // 25: cloud.SearchLogsRequest
(*SearchLogsResponse)(nil), // 26: cloud.SearchLogsResponse
(*SearchLogHit)(nil), // 27: cloud.SearchLogHit
nil, // 28: cloud.InspectContainerResult.LabelsEntry
}
var file_cloud_proto_depIdxs = []int32{
5, // 0: cloud.ToolRequest.list_tools:type_name -> cloud.ListToolsRequest
8, // 1: cloud.ToolRequest.call_tool:type_name -> cloud.CallToolRequest
10, // 2: cloud.ToolRequest.cancel_stream:type_name -> cloud.CancelStreamRequest
6, // 3: cloud.ToolResponse.list_tools:type_name -> cloud.ListToolsResponse
9, // 4: cloud.ToolResponse.call_tool:type_name -> cloud.CallToolResponse
7, // 0: cloud.ToolRequest.list_tools:type_name -> cloud.ListToolsRequest
10, // 1: cloud.ToolRequest.call_tool:type_name -> cloud.CallToolRequest
12, // 2: cloud.ToolRequest.cancel_stream:type_name -> cloud.CancelStreamRequest
8, // 3: cloud.ToolResponse.list_tools:type_name -> cloud.ListToolsResponse
11, // 4: cloud.ToolResponse.call_tool:type_name -> cloud.CallToolResponse
3, // 5: cloud.ToolResponse.log_batch:type_name -> cloud.LogBatch
4, // 6: cloud.LogBatch.entries:type_name -> cloud.LogBatchEntry
7, // 7: cloud.ListToolsResponse.tools:type_name -> cloud.ToolDefinition
0, // 8: cloud.ToolDefinition.scope:type_name -> cloud.ToolScope
12, // 9: cloud.CallToolResponse.list_hosts:type_name -> cloud.ListHostsResult
14, // 10: cloud.CallToolResponse.list_containers:type_name -> cloud.ListContainersResult
16, // 11: cloud.CallToolResponse.container_stats:type_name -> cloud.ContainerStatsResult
20, // 12: cloud.CallToolResponse.action:type_name -> cloud.ActionResult
18, // 13: cloud.CallToolResponse.fetch_logs:type_name -> cloud.FetchLogsResult
19, // 14: cloud.CallToolResponse.inspect_container:type_name -> cloud.InspectContainerResult
21, // 15: cloud.CallToolResponse.deploy:type_name -> cloud.DeployResult
22, // 16: cloud.CallToolResponse.notification:type_name -> cloud.NotificationResult
11, // 17: cloud.ListHostsResult.hosts:type_name -> cloud.HostInfo
13, // 18: cloud.ListContainersResult.containers:type_name -> cloud.ContainerInfo
15, // 19: cloud.ContainerStatsResult.stats:type_name -> cloud.ContainerStatEntry
17, // 20: cloud.FetchLogsResult.entries:type_name -> cloud.LogEntry
26, // 21: cloud.InspectContainerResult.labels:type_name -> cloud.InspectContainerResult.LabelsEntry
25, // 22: cloud.SearchLogsResponse.hits:type_name -> cloud.SearchLogHit
2, // 23: cloud.CloudToolService.ToolStream:input_type -> cloud.ToolResponse
23, // 24: cloud.CloudToolService.SearchLogs:input_type -> cloud.SearchLogsRequest
1, // 25: cloud.CloudToolService.ToolStream:output_type -> cloud.ToolRequest
24, // 26: cloud.CloudToolService.SearchLogs:output_type -> cloud.SearchLogsResponse
25, // [25:27] is the sub-list for method output_type
23, // [23:25] is the sub-list for method input_type
23, // [23:23] is the sub-list for extension type_name
23, // [23:23] is the sub-list for extension extendee
0, // [0:23] is the sub-list for field type_name
5, // 6: cloud.ToolResponse.stats_batch:type_name -> cloud.StatsBatch
4, // 7: cloud.LogBatch.entries:type_name -> cloud.LogBatchEntry
6, // 8: cloud.StatsBatch.entries:type_name -> cloud.StatsBatchEntry
9, // 9: cloud.ListToolsResponse.tools:type_name -> cloud.ToolDefinition
0, // 10: cloud.ToolDefinition.scope:type_name -> cloud.ToolScope
14, // 11: cloud.CallToolResponse.list_hosts:type_name -> cloud.ListHostsResult
16, // 12: cloud.CallToolResponse.list_containers:type_name -> cloud.ListContainersResult
18, // 13: cloud.CallToolResponse.container_stats:type_name -> cloud.ContainerStatsResult
22, // 14: cloud.CallToolResponse.action:type_name -> cloud.ActionResult
20, // 15: cloud.CallToolResponse.fetch_logs:type_name -> cloud.FetchLogsResult
21, // 16: cloud.CallToolResponse.inspect_container:type_name -> cloud.InspectContainerResult
23, // 17: cloud.CallToolResponse.deploy:type_name -> cloud.DeployResult
24, // 18: cloud.CallToolResponse.notification:type_name -> cloud.NotificationResult
13, // 19: cloud.ListHostsResult.hosts:type_name -> cloud.HostInfo
15, // 20: cloud.ListContainersResult.containers:type_name -> cloud.ContainerInfo
17, // 21: cloud.ContainerStatsResult.stats:type_name -> cloud.ContainerStatEntry
19, // 22: cloud.FetchLogsResult.entries:type_name -> cloud.LogEntry
28, // 23: cloud.InspectContainerResult.labels:type_name -> cloud.InspectContainerResult.LabelsEntry
27, // 24: cloud.SearchLogsResponse.hits:type_name -> cloud.SearchLogHit
2, // 25: cloud.CloudToolService.ToolStream:input_type -> cloud.ToolResponse
25, // 26: cloud.CloudToolService.SearchLogs:input_type -> cloud.SearchLogsRequest
1, // 27: cloud.CloudToolService.ToolStream:output_type -> cloud.ToolRequest
26, // 28: cloud.CloudToolService.SearchLogs:output_type -> cloud.SearchLogsResponse
27, // [27:29] is the sub-list for method output_type
25, // [25:27] is the sub-list for method input_type
25, // [25:25] is the sub-list for extension type_name
25, // [25:25] is the sub-list for extension extendee
0, // [0:25] is the sub-list for field type_name
}
func init() { file_cloud_proto_init() }
@@ -2456,8 +2713,9 @@ func file_cloud_proto_init() {
(*ToolResponse_ListTools)(nil),
(*ToolResponse_CallTool)(nil),
(*ToolResponse_LogBatch)(nil),
(*ToolResponse_StatsBatch)(nil),
}
file_cloud_proto_msgTypes[8].OneofWrappers = []any{
file_cloud_proto_msgTypes[10].OneofWrappers = []any{
(*CallToolResponse_ListHosts)(nil),
(*CallToolResponse_ListContainers)(nil),
(*CallToolResponse_ContainerStats)(nil),
@@ -2473,7 +2731,7 @@ func file_cloud_proto_init() {
GoPackagePath: reflect.TypeOf(x{}).PkgPath(),
RawDescriptor: unsafe.Slice(unsafe.StringData(file_cloud_proto_rawDesc), len(file_cloud_proto_rawDesc)),
NumEnums: 1,
NumMessages: 26,
NumMessages: 28,
NumExtensions: 0,
NumServices: 1,
},
+48
View File
@@ -33,6 +33,10 @@ message ToolResponse {
// Dozzle to Cloud for ingestion into VictoriaLogs. request_id is empty
// for log batches — they are not replies to a ToolRequest.
LogBatch log_batch = 4;
// Unsolicited server-push: 30s-windowed container resource metrics.
// request_id is empty for stats batches — like log batches, they are not
// replies to a ToolRequest.
StatsBatch stats_batch = 5;
}
}
@@ -59,6 +63,50 @@ message LogBatchEntry {
uint32 log_id = 8;
}
// Batch of windowed container resource metrics. Dozzle aggregates raw ~1Hz
// stats into fixed windows and pushes one batch per window while connected.
// Cloud writes them to VictoriaMetrics scoped to the owning user (derived
// from the connection's auth).
message StatsBatch {
repeated StatsBatchEntry entries = 1;
}
// One container's metrics for one aggregation window.
//
// Series identity on the Cloud side is (host_id, container_name) — there is
// deliberately no container_id. Container ids churn on every redeploy and each
// new value would strand the old time series forever; keying on the name means
// a redeploy continues the same series, which is what a resource chart wants.
message StatsBatchEntry {
string host_id = 1;
string container_name = 2;
int64 timestamp_ns = 3; // window END, unix nanoseconds
// Gauges — arithmetic mean over the window.
double cpu_percent = 4; // normalised 0-100 across cores, NOT per-core
double memory_percent = 5;
double memory_usage_bytes = 6;
// Peak sample in the window. A 30s mean flattens spikes, and spikes are
// exactly what an incident chart is read for.
double cpu_percent_max = 7;
double memory_percent_max = 8;
// Monotonic counters — the LAST value observed in the window, not a delta.
// Cloud stores these as Prometheus counters so PromQL rate()/increase()
// handle the reset that happens when a container is recreated.
uint64 network_rx_total = 9;
uint64 network_tx_total = 10;
uint64 disk_read_total = 11;
uint64 disk_write_total = 12;
// Number of raw samples folded into this entry. Windows with zero samples
// are never emitted.
uint32 samples = 13;
// The core count cpu_percent was divided by, so the raw per-core figure
// Docker reports stays recoverable downstream.
double cpu_cores = 14;
}
message ListToolsRequest {}
message ListToolsResponse {