feat(cloud): report what kind of process this is, and tell agents about streamLogs (#4996)

Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
Amir Raminfar
2026-09-04 21:54:51 +00:00
committed by GitHub
co-authored by Claude Opus 5
parent fed01ebc28
commit c838ebd360
14 changed files with 228 additions and 21 deletions
+73
View File
@@ -1,6 +1,51 @@
#!/bin/bash
set -e
usage() {
cat <<'EOF'
Usage: ./setup-remote-agent.sh [options] [VM_NAME] [DISTRO] [AGENT_PORT]
Creates an OrbStack VM, installs Docker and runs a Dozzle agent in it.
Positional arguments:
VM_NAME OrbStack VM to create (default: dozzle-agent)
DISTRO Distro for the VM (default: ubuntu)
AGENT_PORT Host port for the agent (default: 7007)
Options:
--local Build and load amir20/dozzle:local instead of pulling :latest
--agent-url URL Dozzle Cloud gRPC endpoint the agent connects to
(default: https://agent.doligence.dozzle.dev)
--cloud-url URL Dozzle Cloud HTTP API the agent posts notifications to
(default: https://doligence.dozzle.dev)
-h, --help Show this help
Both URLs can also come from the environment, as AGENT_URL and DOLIGENCE_URL —
the same names the agent itself reads — so a shell that already exports them
needs no flags. Flags win over the environment.
Pointing an agent at a local Dozzle Cloud stack:
./setup-remote-agent.sh --local \
--agent-url http://host.orb.internal:8082 \
--cloud-url http://host.orb.internal:8080
host.orb.internal is the Mac host, and 8082/8080 are the gRPC and public ports
compose.override.yaml publishes there. An http:// URL makes the agent dial gRPC
in plaintext, which is what the local stack serves.
Use the name, not OrbStack's 198.19.249.2 host address: that one is only routable
from a Docker container running under OrbStack directly. The agent this script
sets up runs inside a Linux VM, which reaches the Mac over its own bridge
instead, and dialling 198.19.249.2 from there times out with nothing to explain
why. The name resolves from both.
Setting these does not by itself connect the agent to Dozzle Cloud: it dials
only once it has an API key, which arrives when the main Dozzle instance pushes
its cloud config down and the agent persists it to /data/cloud.yml.
EOF
}
# Parse arguments
USE_LOCAL=false
POSITIONAL_ARGS=()
@@ -11,6 +56,18 @@ while [[ $# -gt 0 ]]; do
USE_LOCAL=true
shift
;;
--agent-url)
AGENT_URL="$2"
shift 2
;;
--cloud-url)
DOLIGENCE_URL="$2"
shift 2
;;
-h|--help)
usage
exit 0
;;
*)
POSITIONAL_ARGS+=("$1")
shift
@@ -26,6 +83,14 @@ SHARED_CERT="./shared_cert.pem"
SHARED_KEY="./shared_key.pem"
DOZZLE_IMAGE="amir20/dozzle:latest"
# Dozzle Cloud endpoints. These are the agent's own env vars, passed through
# rather than reinvented: AGENT_URL is the gRPC endpoint it dials for tool
# dispatch and log streaming, DOLIGENCE_URL is the HTTP API its cloud
# dispatcher posts notifications to. They are separate on purpose — setting
# only one leaves the other half of the agent talking to production.
CLOUD_AGENT_URL="${AGENT_URL:-https://agent.doligence.dozzle.dev}"
CLOUD_API_URL="${DOLIGENCE_URL:-https://doligence.dozzle.dev}"
if [ "$USE_LOCAL" = true ]; then
DOZZLE_IMAGE="amir20/dozzle:local"
fi
@@ -34,6 +99,8 @@ echo "🚀 Setting up Dozzle Agent on OrbStack VM: $VM_NAME"
if [ "$USE_LOCAL" = true ]; then
echo " Using locally built image"
fi
echo " Dozzle Cloud gRPC: $CLOUD_AGENT_URL"
echo " Dozzle Cloud API: $CLOUD_API_URL"
# Verify shared certificates exist
if [ ! -f "$SHARED_CERT" ]; then
@@ -121,6 +188,8 @@ docker run -d --name dozzle-agent \
-v ~/dozzle-data:/data \
-p $AGENT_PORT:7007 \
-e DOZZLE_LEVEL=debug \
-e AGENT_URL=$CLOUD_AGENT_URL \
-e DOLIGENCE_URL=$CLOUD_API_URL \
$DOZZLE_IMAGE agent \
--cert /certs/shared_cert.pem \
--key /certs/shared_key.pem
@@ -157,6 +226,8 @@ echo " docker run -v /var/run/docker.sock:/var/run/docker.sock \\"
echo " -v $PWD/shared_cert.pem:/shared_cert.pem:ro \\"
echo " -v $PWD/shared_key.pem:/shared_key.pem:ro \\"
echo " -p 8080:8080 \\"
echo " -e AGENT_URL=$CLOUD_AGENT_URL \\"
echo " -e DOLIGENCE_URL=$CLOUD_API_URL \\"
echo " amir20/dozzle:latest \\"
echo " --remote-agent $VM_NAME.orb.local:$AGENT_PORT \\"
echo " --cert /shared_cert.pem --key /shared_key.pem"
@@ -166,6 +237,8 @@ echo ""
echo " DOZZLE_REMOTE_AGENT: $VM_NAME.orb.local:$AGENT_PORT"
echo " DOZZLE_CERT: /shared_cert.pem"
echo " DOZZLE_KEY: /shared_key.pem"
echo " AGENT_URL: $CLOUD_AGENT_URL"
echo " DOLIGENCE_URL: $CLOUD_API_URL"
echo ""
echo "Useful commands:"
echo " View agent logs: orb exec -m $VM_NAME docker logs -f dozzle-agent"
+3 -2
View File
@@ -650,8 +650,9 @@ func (c *Client) UpdateCloudConfig(ctx context.Context, cloudConfig *types.Cloud
req := &pb.UpdateCloudConfigRequest{}
if cloudConfig != nil {
req.CloudConfig = &pb.NotificationCloudConfig{
ApiKey: cloudConfig.APIKey,
Prefix: cloudConfig.Prefix,
ApiKey: cloudConfig.APIKey,
Prefix: cloudConfig.Prefix,
StreamLogs: cloudConfig.StreamLogs,
}
if cloudConfig.ExpiresAt != nil {
req.CloudConfig.ExpiresAt = timestamppb.New(*cloudConfig.ExpiresAt)
+1
View File
@@ -38,6 +38,7 @@ func (m *mockNotificationHandler) HandleNotificationConfig(subscriptions []types
}
func (m *mockNotificationHandler) SetCloudDispatcher(d dispatcher.Dispatcher) {}
func (m *mockNotificationHandler) SetCloudStreamLogs(enabled *bool) {}
func (m *mockNotificationHandler) ClearCloudDispatcher() {}
func (m *mockNotificationHandler) GetNotificationStats() []types.SubscriptionStats {
+24 -7
View File
@@ -1,7 +1,7 @@
// Code generated by protoc-gen-go. DO NOT EDIT.
// versions:
// protoc-gen-go v1.36.11
// protoc v7.35.1
// protoc v7.36.0
// source: types.proto
package pb
@@ -1264,10 +1264,15 @@ func (x *NotificationDispatcher) GetHeaders() map[string]string {
}
type NotificationCloudConfig struct {
state protoimpl.MessageState `protogen:"open.v1"`
ApiKey string `protobuf:"bytes,1,opt,name=apiKey,proto3" json:"apiKey,omitempty"`
Prefix string `protobuf:"bytes,2,opt,name=prefix,proto3" json:"prefix,omitempty"`
ExpiresAt *timestamppb.Timestamp `protobuf:"bytes,3,opt,name=expiresAt,proto3" json:"expiresAt,omitempty"`
state protoimpl.MessageState `protogen:"open.v1"`
ApiKey string `protobuf:"bytes,1,opt,name=apiKey,proto3" json:"apiKey,omitempty"`
Prefix string `protobuf:"bytes,2,opt,name=prefix,proto3" json:"prefix,omitempty"`
ExpiresAt *timestamppb.Timestamp `protobuf:"bytes,3,opt,name=expiresAt,proto3" json:"expiresAt,omitempty"`
// Whether the fleet streams container logs to Dozzle Cloud. Optional so an
// agent can tell "the hub did not say" (an older hub — keep the historical
// default of enabled) from "the hub said no". Without this an agent never
// learned the setting at all and streamed regardless of what the user chose.
StreamLogs *bool `protobuf:"varint,4,opt,name=streamLogs,proto3,oneof" json:"streamLogs,omitempty"`
unknownFields protoimpl.UnknownFields
sizeCache protoimpl.SizeCache
}
@@ -1323,6 +1328,13 @@ func (x *NotificationCloudConfig) GetExpiresAt() *timestamppb.Timestamp {
return nil
}
func (x *NotificationCloudConfig) GetStreamLogs() bool {
if x != nil && x.StreamLogs != nil {
return *x.StreamLogs
}
return false
}
type NotificationSubscriptionStats struct {
state protoimpl.MessageState `protogen:"open.v1"`
SubscriptionId int32 `protobuf:"varint,1,opt,name=subscriptionId,proto3" json:"subscriptionId,omitempty"`
@@ -1520,11 +1532,15 @@ const file_types_proto_rawDesc = "" +
"\fHeadersEntry\x12\x10\n" +
"\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n" +
"\x05value\x18\x02 \x01(\tR\x05value:\x028\x01J\x04\b\a\x10\bJ\x04\b\b\x10\tJ\x04\b\t\x10\n" +
"\"\x83\x01\n" +
"\"\xb7\x01\n" +
"\x17NotificationCloudConfig\x12\x16\n" +
"\x06apiKey\x18\x01 \x01(\tR\x06apiKey\x12\x16\n" +
"\x06prefix\x18\x02 \x01(\tR\x06prefix\x128\n" +
"\texpiresAt\x18\x03 \x01(\v2\x1a.google.protobuf.TimestampR\texpiresAt\"\xe7\x01\n" +
"\texpiresAt\x18\x03 \x01(\v2\x1a.google.protobuf.TimestampR\texpiresAt\x12#\n" +
"\n" +
"streamLogs\x18\x04 \x01(\bH\x00R\n" +
"streamLogs\x88\x01\x01B\r\n" +
"\v_streamLogs\"\xe7\x01\n" +
"\x1dNotificationSubscriptionStats\x12&\n" +
"\x0esubscriptionId\x18\x01 \x01(\x05R\x0esubscriptionId\x12\"\n" +
"\ftriggerCount\x18\x02 \x01(\x03R\ftriggerCount\x12D\n" +
@@ -1606,6 +1622,7 @@ func file_types_proto_init() {
if File_types_proto != nil {
return
}
file_types_proto_msgTypes[13].OneofWrappers = []any{}
type x struct{}
out := protoimpl.TypeBuilder{
File: protoimpl.DescBuilder{
+6
View File
@@ -33,6 +33,9 @@ import (
type NotificationConfigHandler interface {
HandleNotificationConfig(subscriptions []types.SubscriptionConfig, dispatchers []types.DispatcherConfig) error
SetCloudDispatcher(d dispatcher.Dispatcher)
// SetCloudStreamLogs applies the hub's log-streaming choice. nil means the
// hub did not say (an older build), which keeps the historical default.
SetCloudStreamLogs(enabled *bool)
ClearCloudDispatcher()
GetNotificationStats() []types.SubscriptionStats
}
@@ -523,6 +526,9 @@ func (s *server) UpdateCloudConfig(ctx context.Context, req *pb.UpdateCloudConfi
return nil, status.Error(codes.Internal, err.Error())
}
s.notificationConfigHandler.SetCloudDispatcher(d)
// After SetCloudDispatcher, which rebuilds the stored config from the
// dispatcher and so cannot carry this itself.
s.notificationConfigHandler.SetCloudStreamLogs(cc.StreamLogs)
log.Info().Msg("Updated cloud config from main server")
} else {
s.notificationConfigHandler.ClearCloudDispatcher()
+29 -4
View File
@@ -35,10 +35,17 @@ const (
// Client manages the gRPC connection to Dozzle Cloud
type Client struct {
deps ToolDeps
apiKeyFunc func() string
instanceID string
version string
deps ToolDeps
apiKeyFunc func() string
instanceID string
version string
// mode and swarmClusterID describe what kind of Dozzle process this is.
// Cloud cannot infer either — a swarm replica and a standalone server both
// scope themselves to their own daemon and look identical on the wire — so
// they travel as connect-time metadata. Both may be empty; cloud treats
// that as unknown.
mode string
swarmClusterID string
streamLogsFunc func() bool
target string
plaintext bool
@@ -102,6 +109,18 @@ func NewClient(apiKeyFunc func() string, instanceID string, version string, deps
}
}
// SetDeployment records what kind of Dozzle process this is ("server",
// "swarm", "k8s" or "agent") and, on a swarm node, which swarm it belongs to.
// Sent as connect-time metadata so Dozzle Cloud can tell a hub from the agents
// behind it, and a replicated swarm from one API key reused across deployments.
//
// Optional: a client that never calls this simply reports nothing, which is
// what every Dozzle before this change did.
func (c *Client) SetDeployment(mode, swarmClusterID string) {
c.mode = mode
c.swarmClusterID = swarmClusterID
}
// SetStreamLogsFunc registers a function that reports whether bulk container
// log streaming to cloud is enabled. If unset, the streamer runs by default.
func (c *Client) SetStreamLogsFunc(f func() bool) {
@@ -254,6 +273,12 @@ func (c *Client) connect(ctx context.Context, apiKey string) (wasConnected bool,
if c.instanceID != "" {
mdPairs = append(mdPairs, "x-instance-id", c.instanceID)
}
if c.mode != "" {
mdPairs = append(mdPairs, "x-dozzle-mode", c.mode)
}
if c.swarmClusterID != "" {
mdPairs = append(mdPairs, "x-dozzle-swarm-cluster-id", c.swarmClusterID)
}
md := metadata.Pairs(mdPairs...)
streamCtx := metadata.NewOutgoingContext(connCtx, md)
+25
View File
@@ -0,0 +1,25 @@
package cloud
import "testing"
// Cloud cannot infer what kind of process is connecting — a swarm replica and a
// standalone server both scope themselves to their own daemon and look
// identical on the wire — so Dozzle reports it. Every build before this one
// reported nothing, and those connections must keep working.
func TestSetDeployment(t *testing.T) {
c := &Client{}
if c.mode != "" || c.swarmClusterID != "" {
t.Fatal("a client that never calls SetDeployment must report nothing")
}
c.SetDeployment("swarm", "cluster-abc")
if c.mode != "swarm" || c.swarmClusterID != "cluster-abc" {
t.Fatalf("got mode=%q cluster=%q", c.mode, c.swarmClusterID)
}
// Outside a swarm there is no cluster to report, but the mode still stands.
c.SetDeployment("server", "")
if c.mode != "server" || c.swarmClusterID != "" {
t.Fatalf("got mode=%q cluster=%q", c.mode, c.swarmClusterID)
}
}
+6 -1
View File
@@ -27,7 +27,12 @@ type Host struct {
Type string `json:"type"`
Available bool `json:"available"`
Swarm bool `json:"-"`
Group string `json:"group,omitempty"`
// SwarmClusterID identifies the swarm this node belongs to, empty outside
// a swarm. Every node of one swarm reports the same value, which is what
// lets Dozzle Cloud tell "one swarm, N replicas" apart from one API key
// reused across unrelated deployments.
SwarmClusterID string `json:"-"`
Group string `json:"group,omitempty"`
}
func (h Host) String() string {
+3
View File
@@ -70,6 +70,9 @@ func NewClient(cli DockerCLI, host container.Host) *DockerClient {
host.DockerVersion = info.ServerVersion
host.Runtime = detectRuntime(cli, info)
host.Swarm = info.Swarm.NodeID != ""
if info.Swarm.Cluster != nil {
host.SwarmClusterID = info.Swarm.Cluster.ID
}
return &DockerClient{
cli: cli,
+24 -1
View File
@@ -72,6 +72,21 @@ func (h *persistingNotificationHandler) HandleNotificationConfig(subscriptions [
return nil
}
// SetCloudStreamLogs applies the hub's log-streaming choice to this agent's own
// cloud client. The agent streams its logs to cloud itself rather than through
// the hub, so without this a user who turned streaming off saw it stop on the
// hub while every agent kept sending.
func (h *persistingNotificationHandler) SetCloudStreamLogs(enabled *bool) {
cc := h.cloudConfig.Load()
if cc == nil {
return
}
updated := *cc
updated.StreamLogs = enabled
h.cloudConfig.Store(&updated)
h.persistCloudConfig(updated)
}
func (h *persistingNotificationHandler) SetCloudDispatcher(d dispatcher.Dispatcher) {
h.manager.SetCloudDispatcher(d)
@@ -90,6 +105,11 @@ func (h *persistingNotificationHandler) SetCloudDispatcher(d dispatcher.Dispatch
if h.onCloudSet != nil {
h.onCloudSet()
}
h.persistCloudConfig(cc)
}
// persistCloudConfig writes cloud.yml so the setting survives an agent restart.
func (h *persistingNotificationHandler) persistCloudConfig(cc notification.CloudConfig) {
if err := os.MkdirAll("./data", 0755); err != nil {
log.Error().Err(err).Msg("Could not create data directory for cloud config")
return
@@ -203,9 +223,10 @@ func (a *AgentCmd) Run(args Args, embeddedCerts embed.FS) error {
// Cloud gRPC client — connects directly to Dozzle Cloud with this agent's
// own host ID as instance_id, so log streaming and tool dispatch happen
// here instead of funneling through the main server.
var instanceID string
var instanceID, swarmClusterID string
if h, err := agentHostService.LocalHost(); err == nil {
instanceID = h.ID
swarmClusterID = h.SwarmClusterID
}
apiKeyFunc := func() string {
if cc := notificationHandler.CloudConfig(); cc != nil {
@@ -218,6 +239,8 @@ func (a *AgentCmd) Run(args Args, embeddedCerts embed.FS) error {
HostService: agentHostService,
Labels: args.Filter,
})
// An agent is always an agent, whatever the hub in front of it is running as.
cloudClient.SetDeployment("agent", swarmClusterID)
cloudClient.SetStreamLogsFunc(func() bool {
cc := notificationHandler.CloudConfig()
return cc != nil && cc.StreamLogsEnabled()
+21 -5
View File
@@ -274,10 +274,14 @@ func (m *MultiHostService) SetCloudConfig(cc *notification.CloudConfig) {
}
// SetCloudStreamLogs updates the bulk-log-streaming privacy flag on the cloud
// config and persists it. Only affects the local cloud client; agents never
// stream logs directly to cloud.
// config, persists it, and broadcasts it.
//
// Agents connect to cloud themselves and stream their own logs, so the setting
// has to reach them: without the broadcast a user who turned streaming off saw
// it stop on the hub while every agent kept sending.
func (m *MultiHostService) SetCloudStreamLogs(enabled bool) {
m.persister.SetCloudStreamLogs(enabled)
m.broadcastCloudConfig()
}
// ResetCloudDispatcherBreaker clears the cloud dispatcher's auth circuit breaker
@@ -358,9 +362,10 @@ func (m *MultiHostService) broadcastCloudConfig() {
var cc *types.CloudConfig
if ncc != nil {
cc = &types.CloudConfig{
APIKey: ncc.APIKey,
Prefix: ncc.Prefix,
ExpiresAt: ncc.ExpiresAt,
APIKey: ncc.APIKey,
Prefix: ncc.Prefix,
ExpiresAt: ncc.ExpiresAt,
StreamLogs: ncc.StreamLogs,
}
}
@@ -443,6 +448,17 @@ func (h *swarmNotificationHandler) SetCloudDispatcher(d dispatcher.Dispatcher) {
h.notify()
}
// SetCloudStreamLogs applies a peer replica's log-streaming choice. Routed
// through the persister like everything else here so disk and manager stay in
// lockstep across replicas.
func (h *swarmNotificationHandler) SetCloudStreamLogs(enabled *bool) {
if enabled == nil {
return
}
h.persister.SetCloudStreamLogs(*enabled)
h.notify()
}
func (h *swarmNotificationHandler) ClearCloudDispatcher() {
h.persister.RemoveCloudConfig()
h.notify()
+3 -1
View File
@@ -143,9 +143,10 @@ func main() {
return ""
}
var instanceID string
var instanceID, swarmClusterID string
if h, err := hostService.LocalHost(); err == nil {
instanceID = h.ID
swarmClusterID = h.SwarmClusterID
}
cloudHostService := newCloudHostService(args.Mode, hostService)
@@ -156,6 +157,7 @@ func main() {
Labels: args.Filter,
NotificationService: notificationService,
})
cloudClient.SetDeployment(args.Mode, swarmClusterID)
cloudClient.SetStreamLogsFunc(func() bool {
return hostService.CloudConfig().StreamLogsEnabled()
})
+5
View File
@@ -153,6 +153,11 @@ message NotificationCloudConfig {
string apiKey = 1;
string prefix = 2;
google.protobuf.Timestamp expiresAt = 3;
// Whether the fleet streams container logs to Dozzle Cloud. Optional so an
// agent can tell "the hub did not say" (an older hub — keep the historical
// default of enabled) from "the hub said no". Without this an agent never
// learned the setting at all and streamed regardless of what the user chose.
optional bool streamLogs = 4;
}
message NotificationSubscriptionStats {
+5
View File
@@ -103,6 +103,11 @@ type CloudConfig struct {
APIKey string
Prefix string
ExpiresAt *time.Time
// StreamLogs is the user's log-streaming choice, nil when the sender did
// not say. Agents apply it to their own cloud client; nil keeps the
// historical default (enabled) so an older hub doesn't silently turn
// streaming off.
StreamLogs *bool
}
// DispatcherConfig represents a dispatcher configuration