From 2b92500c01decd992361bdcd77653e84188a059e Mon Sep 17 00:00:00 2001 From: Thibault Le Ouay Ducasse Date: Tue, 21 Jul 2026 18:01:46 +0200 Subject: [PATCH] private location: typo and bug refresh --- apps/checker/pkg/scheduler/scheduler.go | 333 ++++++++++-------- apps/checker/pkg/scheduler/scheduler_test.go | 80 +++++ .../pages/docs/concept/private-locations.mdx | 2 +- .../docs/concept/probes-and-locations.mdx | 2 +- .../guides/how-to-create-private-location.mdx | 6 +- 5 files changed, 269 insertions(+), 154 deletions(-) diff --git a/apps/checker/pkg/scheduler/scheduler.go b/apps/checker/pkg/scheduler/scheduler.go index 7b5e3860..68c4413f 100644 --- a/apps/checker/pkg/scheduler/scheduler.go +++ b/apps/checker/pkg/scheduler/scheduler.go @@ -1,6 +1,7 @@ package scheduler import ( + "bytes" "context" "log" "sync" @@ -10,6 +11,7 @@ import ( "github.com/madflojo/tasks" "github.com/openstatushq/openstatus/apps/checker/pkg/job" v1 "github.com/openstatushq/openstatus/apps/checker/proto/private_location/v1" + "google.golang.org/protobuf/proto" ) const ( @@ -27,6 +29,40 @@ type MonitorManager struct { JobRunner job.JobRunner Scheduler *tasks.Scheduler mu sync.Mutex + configs map[string][]byte +} + +// shouldSchedule reports whether a task has to be created for the monitor, and +// drops the running one first when the config changed: a task captures its +// monitor when it is created, so an edited monitor would otherwise keep +// checking with the config it had on the probe's first fetch. +func (mm *MonitorManager) shouldSchedule(id string, monitor proto.Message) bool { + mm.mu.Lock() + defer mm.mu.Unlock() + + if mm.configs == nil { + mm.configs = make(map[string][]byte) + } + + config, err := proto.MarshalOptions{Deterministic: true}.Marshal(monitor) + if err != nil { + log.Printf("Failed to encode config for monitor %s: %v", id, err) + config = nil + } + + if _, lookupErr := mm.Scheduler.Lookup(id); lookupErr != nil { + mm.configs[id] = config + return true + } + + if bytes.Equal(mm.configs[id], config) { + return false + } + + log.Printf("Config changed for monitor %s, rescheduling", id) + mm.Scheduler.Del(id) + mm.configs[id] = config + return true } // UpdateMonitors fetches the latest monitors and starts/stops jobs as needed @@ -42,176 +78,175 @@ func (mm *MonitorManager) UpdateMonitors(ctx context.Context) { // HTTP monitors: start jobs for new monitors for _, m := range res.Msg.HttpMonitors { currentIDs[m.Id] = struct{}{} - _, err := mm.Scheduler.Lookup(m.Id) - if err != nil { - interval := time.Duration(intervalToSecond(m.Periodicity)) * time.Second - task := tasks.Task{ - Interval: interval, - RunOnce: false, - RunSingleInstance: true, - // StartAfter: time.Duration(1) * time.Second, - ErrFunc: func(e error) { - log.Printf("An error occurred when executing task %s", e) - }, - FuncWithTaskContext: func(ctx tasks.TaskContext) error { - monitor := m - c := context.Background() - log.Printf("Starting job for monitor %s (%s)", monitor.Id, monitor.Url) - data, err := mm.JobRunner.HTTPJob(c, monitor, res.Msg.Region) - - if err != nil { - log.Printf("Monitor check failed for %s (%s): %v", monitor.Id, monitor.Url, err) - return err - } - resp, ingestErr := mm.Client.IngestHTTP(c, &connect.Request[v1.IngestHTTPRequest]{ - Msg: &v1.IngestHTTPRequest{ - MonitorId: monitor.Id, - Id: data.ID, - Url: monitor.Url, - Message: data.Message, - Latency: data.Latency, - Timing: data.Timing, - Headers: data.Headers, - Body: data.Body, - RequestStatus: data.RequestStatus, - StatusCode: int64(data.StatusCode), - Error: int64(data.Error), - CronTimestamp: data.CronTimestamp, - Timestamp: data.Timestamp, - }, - }) - if ingestErr != nil { - log.Printf("Failed to ingest HTTP result for %s (%s): %v", monitor.Id, monitor.Url, ingestErr) - return ingestErr - } - log.Printf("Monitor check succeeded for %s (%s), ingest response: %v", monitor.Id, monitor.Url, resp) - return nil - }, - } - - err := mm.Scheduler.AddWithID(m.Id, &task) - - if err != nil { - log.Printf("Failed to add HTTP monitor job for %s (%s): %v", m.Id, m.Url, err) - continue - } - log.Printf("Started monitoring job for %s (%s)", m.Id, m.Url) + if !mm.shouldSchedule(m.Id, m) { continue } + interval := time.Duration(intervalToSecond(m.Periodicity)) * time.Second + task := tasks.Task{ + Interval: interval, + RunOnce: false, + RunSingleInstance: true, + // StartAfter: time.Duration(1) * time.Second, + ErrFunc: func(e error) { + log.Printf("An error occurred when executing task %s", e) + }, + FuncWithTaskContext: func(ctx tasks.TaskContext) error { + monitor := m + c := context.Background() + log.Printf("Starting job for monitor %s (%s)", monitor.Id, monitor.Url) + data, err := mm.JobRunner.HTTPJob(c, monitor, res.Msg.Region) + + if err != nil { + log.Printf("Monitor check failed for %s (%s): %v", monitor.Id, monitor.Url, err) + return err + } + resp, ingestErr := mm.Client.IngestHTTP(c, &connect.Request[v1.IngestHTTPRequest]{ + Msg: &v1.IngestHTTPRequest{ + MonitorId: monitor.Id, + Id: data.ID, + Url: monitor.Url, + Message: data.Message, + Latency: data.Latency, + Timing: data.Timing, + Headers: data.Headers, + Body: data.Body, + RequestStatus: data.RequestStatus, + StatusCode: int64(data.StatusCode), + Error: int64(data.Error), + CronTimestamp: data.CronTimestamp, + Timestamp: data.Timestamp, + }, + }) + if ingestErr != nil { + log.Printf("Failed to ingest HTTP result for %s (%s): %v", monitor.Id, monitor.Url, ingestErr) + return ingestErr + } + log.Printf("Monitor check for %s (%s) ingested with status %q (code %d), ingest response: %v", monitor.Id, monitor.Url, data.RequestStatus, data.StatusCode, resp) + return nil + }, + } + + if err := mm.Scheduler.AddWithID(m.Id, &task); err != nil { + log.Printf("Failed to add HTTP monitor job for %s (%s): %v", m.Id, m.Url, err) + continue + } + log.Printf("Started monitoring job for %s (%s)", m.Id, m.Url) } // TCP monitors: start jobs for new monitors for _, m := range res.Msg.TcpMonitors { currentIDs[m.Id] = struct{}{} - _, err := mm.Scheduler.Lookup(m.Id) - if err != nil { - - interval := time.Duration(intervalToSecond(m.Periodicity)) * time.Second - task := tasks.Task{ - Interval: interval, - RunOnce: false, - // StartAfter: time.Now().Add(5 * time.Millisecond), - RunSingleInstance: true, - FuncWithTaskContext: func(ctx tasks.TaskContext) error { - - monitor := m - c := context.Background() - log.Printf("Starting TCP job for monitor %s (%s)", monitor.Id, monitor.Uri) - data, err := mm.JobRunner.TCPJob(c, monitor, res.Msg.Region) - if err != nil { - log.Printf("TCP monitor check failed for %s (%s): %v", monitor.Id, monitor.Uri, err) - } - resp, ingestErr := mm.Client.IngestTCP(c, &connect.Request[v1.IngestTCPRequest]{ - Msg: &v1.IngestTCPRequest{ - MonitorId: monitor.Id, - Id: data.ID, - Uri: monitor.Uri, - Message: data.Message, - Latency: data.Latency, - RequestStatus: data.RequestStatus, - Error: int64(data.Error), - CronTimestamp: data.CronTimestamp, - Timestamp: data.Timestamp, - }, - }) - if ingestErr != nil { - log.Printf("Failed to ingest TCP result for %s (%s): %v", monitor.Id, monitor.Uri, ingestErr) - return ingestErr - } - log.Printf("TCP monitor check succeeded for %s (%s), ingest response: %v", monitor.Id, monitor.Uri, resp) - - return nil - }, - } - err := mm.Scheduler.AddWithID(m.Id, &task) - - if err != nil { - log.Printf("Failed to add TCP monitor job for %s (%s): %v", m.Id, m.Uri, err) - continue - } - log.Printf("Started TCP monitoring job for %s (%s)", m.Id, m.Uri) + if !mm.shouldSchedule(m.Id, m) { + continue + } + + interval := time.Duration(intervalToSecond(m.Periodicity)) * time.Second + task := tasks.Task{ + Interval: interval, + RunOnce: false, + // StartAfter: time.Now().Add(5 * time.Millisecond), + RunSingleInstance: true, + FuncWithTaskContext: func(ctx tasks.TaskContext) error { + + monitor := m + c := context.Background() + log.Printf("Starting TCP job for monitor %s (%s)", monitor.Id, monitor.Uri) + data, err := mm.JobRunner.TCPJob(c, monitor, res.Msg.Region) + if err != nil { + log.Printf("TCP monitor check failed for %s (%s): %v", monitor.Id, monitor.Uri, err) + return err + } + resp, ingestErr := mm.Client.IngestTCP(c, &connect.Request[v1.IngestTCPRequest]{ + Msg: &v1.IngestTCPRequest{ + MonitorId: monitor.Id, + Id: data.ID, + Uri: monitor.Uri, + Message: data.Message, + Latency: data.Latency, + RequestStatus: data.RequestStatus, + Error: int64(data.Error), + CronTimestamp: data.CronTimestamp, + Timestamp: data.Timestamp, + }, + }) + if ingestErr != nil { + log.Printf("Failed to ingest TCP result for %s (%s): %v", monitor.Id, monitor.Uri, ingestErr) + return ingestErr + } + log.Printf("TCP monitor check for %s (%s) ingested with status %q, ingest response: %v", monitor.Id, monitor.Uri, data.RequestStatus, resp) + + return nil + }, } + + if err := mm.Scheduler.AddWithID(m.Id, &task); err != nil { + log.Printf("Failed to add TCP monitor job for %s (%s): %v", m.Id, m.Uri, err) + continue + } + log.Printf("Started TCP monitoring job for %s (%s)", m.Id, m.Uri) } for _, m := range res.Msg.DnsMonitors { currentIDs[m.Id] = struct{}{} - _, err := mm.Scheduler.Lookup(m.Id) - if err != nil { - - interval := time.Duration(intervalToSecond(m.Periodicity)) * time.Second - task := tasks.Task{ - Interval: interval, - RunOnce: false, - // StartAfter: time.Now().Add(5 * time.Millisecond), - RunSingleInstance: true, - FuncWithTaskContext: func(ctx tasks.TaskContext) error { - - monitor := m - c := context.Background() - log.Printf("Starting TCP job for monitor %s (%s)", monitor.Id, monitor.Uri) - _, err := mm.JobRunner.DNSJob(c, monitor) - if err != nil { - log.Printf("TCP monitor check failed for %s (%s): %v", monitor.Id, monitor.Uri, err) - } - resp, ingestErr := mm.Client.IngestDNS(c, &connect.Request[v1.IngestDNSRequest]{ - Msg: &v1.IngestDNSRequest{ - MonitorId: monitor.Id, - - // Id: data.ID, - // Uri: monitor.Uri, - // Message: data.Message, - // Latency: data.Latency, - // RequestStatus: data.RequestStatus, - // Error: int64(data.Error), - // CronTimestamp: data.CronTimestamp, - // Timestamp: data.Timestamp, - }, - }) - if ingestErr != nil { - log.Printf("Failed to ingest TCP result for %s (%s): %v", monitor.Id, monitor.Uri, ingestErr) - return ingestErr - } - log.Printf("TCP monitor check succeeded for %s (%s), ingest response: %v", monitor.Id, monitor.Uri, resp) - - return nil - }, - } - err := mm.Scheduler.AddWithID(m.Id, &task) - - if err != nil { - log.Printf("Failed to add TCP monitor job for %s (%s): %v", m.Id, m.Uri, err) - continue - } - log.Printf("Started TCP monitoring job for %s (%s)", m.Id, m.Uri) + if !mm.shouldSchedule(m.Id, m) { + continue + } + + interval := time.Duration(intervalToSecond(m.Periodicity)) * time.Second + task := tasks.Task{ + Interval: interval, + RunOnce: false, + // StartAfter: time.Now().Add(5 * time.Millisecond), + RunSingleInstance: true, + FuncWithTaskContext: func(ctx tasks.TaskContext) error { + + monitor := m + c := context.Background() + log.Printf("Starting DNS job for monitor %s (%s)", monitor.Id, monitor.Uri) + _, err := mm.JobRunner.DNSJob(c, monitor) + if err != nil { + log.Printf("DNS monitor check failed for %s (%s): %v", monitor.Id, monitor.Uri, err) + } + resp, ingestErr := mm.Client.IngestDNS(c, &connect.Request[v1.IngestDNSRequest]{ + Msg: &v1.IngestDNSRequest{ + MonitorId: monitor.Id, + + // Id: data.ID, + // Uri: monitor.Uri, + // Message: data.Message, + // Latency: data.Latency, + // RequestStatus: data.RequestStatus, + // Error: int64(data.Error), + // CronTimestamp: data.CronTimestamp, + // Timestamp: data.Timestamp, + }, + }) + if ingestErr != nil { + log.Printf("Failed to ingest DNS result for %s (%s): %v", monitor.Id, monitor.Uri, ingestErr) + return ingestErr + } + log.Printf("DNS monitor check for %s (%s) ingested, ingest response: %v", monitor.Id, monitor.Uri, resp) + + return nil + }, + } + + if err := mm.Scheduler.AddWithID(m.Id, &task); err != nil { + log.Printf("Failed to add DNS monitor job for %s (%s): %v", m.Id, m.Uri, err) + continue } + log.Printf("Started DNS monitoring job for %s (%s)", m.Id, m.Uri) } + mm.mu.Lock() for id := range mm.Scheduler.Tasks() { if _, stillExists := currentIDs[id]; !stillExists { mm.Scheduler.Del(id) + delete(mm.configs, id) } } + mm.mu.Unlock() } diff --git a/apps/checker/pkg/scheduler/scheduler_test.go b/apps/checker/pkg/scheduler/scheduler_test.go index e70aeed2..0d86d315 100644 --- a/apps/checker/pkg/scheduler/scheduler_test.go +++ b/apps/checker/pkg/scheduler/scheduler_test.go @@ -23,15 +23,23 @@ type mockJobRunner struct { mu sync.Mutex httpRegion string tcpRegion string + httpMonitor *v1.HTTPMonitor } func (m *mockJobRunner) HTTPJob(ctx context.Context, monitor *v1.HTTPMonitor, region string) (*job.HttpPrivateRegionData, error) { m.HTTPJobCalled.Store(true) m.mu.Lock() m.httpRegion = region + m.httpMonitor = monitor m.mu.Unlock() return &job.HttpPrivateRegionData{}, nil } + +func (m *mockJobRunner) HTTPMonitor() *v1.HTTPMonitor { + m.mu.Lock() + defer m.mu.Unlock() + return m.httpMonitor +} func (m *mockJobRunner) TCPJob(ctx context.Context, monitor *v1.TCPMonitor, region string) (*job.TCPPrivateRegionData, error) { m.TCPJobCalled.Store(true) @@ -158,3 +166,75 @@ func TestMonitorManager_StartAndStopJobs_WithJobRunner(t *testing.T) { } } + +// runScheduledTask executes an already scheduled task synchronously, so a test +// can observe the monitor config its closure captured without waiting for the +// interval to elapse. +func runScheduledTask(t *testing.T, s *tasks.Scheduler, id string) { + t.Helper() + + task, err := s.Lookup(id) + if err != nil { + t.Fatalf("expected a task scheduled for %s: %v", id, err) + } + if err := task.FuncWithTaskContext(tasks.TaskContext{}); err != nil { + t.Fatalf("task %s returned an error: %v", id, err) + } +} + +func TestMonitorManager_ReschedulesOnConfigChange(t *testing.T) { + ctx := t.Context() + + // Long periodicity: the task only runs when the test invokes it. + withoutHeader := &v1.HTTPMonitor{Id: "http1", Url: "https://openstat.us", Periodicity: "1h"} + unchanged := &v1.HTTPMonitor{Id: "http1", Url: "https://openstat.us", Periodicity: "1h"} + withHeader := &v1.HTTPMonitor{ + Id: "http1", Url: "https://openstat.us", Periodicity: "1h", + Headers: []*v1.Headers{{Key: "X-Auth-Token", Value: "secret"}}, + } + + current := withoutHeader + client := &mockClient{ + MonitorsFunc: func(ctx context.Context, req *connect.Request[v1.MonitorsRequest]) (*connect.Response[v1.MonitorsResponse], error) { + return connect.NewResponse(&v1.MonitorsResponse{ + HttpMonitors: []*v1.HTTPMonitor{current}, + Region: "frankfurt-dc1", + }), nil + }, + IngestHTTPFunc: func(ctx context.Context, req *connect.Request[v1.IngestHTTPRequest]) (*connect.Response[v1.IngestHTTPResponse], error) { + return connect.NewResponse(&v1.IngestHTTPResponse{}), nil + }, + } + jobRunner := &mockJobRunner{} + + s := tasks.New() + defer s.Stop() + + mm := &scheduler.MonitorManager{Client: client, JobRunner: jobRunner, Scheduler: s} + + mm.UpdateMonitors(ctx) + runScheduledTask(t, mm.Scheduler, "http1") + if got := jobRunner.HTTPMonitor(); got != withoutHeader { + t.Fatalf("expected the job to run with the fetched monitor, got %v", got) + } + + // Same config, fresh pointer: rescheduling here would reset the interval timer. + current = unchanged + mm.UpdateMonitors(ctx) + runScheduledTask(t, mm.Scheduler, "http1") + if got := jobRunner.HTTPMonitor(); got != withoutHeader { + t.Errorf("expected an unchanged monitor to keep its task, got a rescheduled one") + } + + current = withHeader + mm.UpdateMonitors(ctx) + runScheduledTask(t, mm.Scheduler, "http1") + + got := jobRunner.HTTPMonitor() + if got != withHeader { + t.Fatalf("expected the job to run with the updated monitor, got %v", got) + } + if len(got.Headers) != 1 || got.Headers[0].Key != "X-Auth-Token" || got.Headers[0].Value != "secret" { + t.Errorf("expected the added header to reach the job, got %v", got.Headers) + } +} diff --git a/apps/web/src/content/pages/docs/concept/private-locations.mdx b/apps/web/src/content/pages/docs/concept/private-locations.mdx index 6b02a600..a1f955ac 100644 --- a/apps/web/src/content/pages/docs/concept/private-locations.mdx +++ b/apps/web/src/content/pages/docs/concept/private-locations.mdx @@ -25,7 +25,7 @@ If your endpoint is already public and you only want geographic coverage, the [p A private location is not an inbound port openstatus calls into. The probe initiates everything: -1. You deploy the openstatus probe container — typically `ghcr.io/openstatushq/probe:latest` — somewhere inside the network that needs to be monitored. +1. You deploy the openstatus probe container — typically `ghcr.io/openstatushq/private-location:latest` — somewhere inside the network that needs to be monitored. 2. The probe opens an outbound, authenticated connection to the openstatus platform using the token issued when you created the location. 3. openstatus dispatches the checks assigned to that location down the connection. The probe runs them against your internal endpoints and streams the results back. diff --git a/apps/web/src/content/pages/docs/concept/probes-and-locations.mdx b/apps/web/src/content/pages/docs/concept/probes-and-locations.mdx index 5617c997..0835cfec 100644 --- a/apps/web/src/content/pages/docs/concept/probes-and-locations.mdx +++ b/apps/web/src/content/pages/docs/concept/probes-and-locations.mdx @@ -17,7 +17,7 @@ That's the relationship in shorthand. The rest of this page unpacks each term an A **probe** is the software that actually performs a check. It opens the TCP connection, sends the HTTP request, resolves the DNS record, measures the timing — and reports the result back. There are two kinds: - **Public probes.** Operated by openstatus, deployed across Fly.io, Railway, and Koyeb regions. You don't manage them. -- **Private probes.** A container image (`ghcr.io/openstatushq/probe:latest`) you run on your own infrastructure. See [Understanding Private Locations](/docs/concept/private-locations). +- **Private probes.** A container image (`ghcr.io/openstatushq/private-location:latest`) you run on your own infrastructure. See [Understanding Private Locations](/docs/concept/private-locations). In day-to-day usage you rarely interact with the "probe" abstraction directly — you assign a *location* to a monitor, and the location decides which probe runs the check. The word matters mostly when you're reading debug output, opening a support ticket, or deploying a private location. diff --git a/apps/web/src/content/pages/docs/guides/how-to-create-private-location.mdx b/apps/web/src/content/pages/docs/guides/how-to-create-private-location.mdx index 3181a942..30d068f8 100644 --- a/apps/web/src/content/pages/docs/guides/how-to-create-private-location.mdx +++ b/apps/web/src/content/pages/docs/guides/how-to-create-private-location.mdx @@ -37,7 +37,7 @@ docker run -d \ --name openstatus-probe \ --restart unless-stopped \ -e OPENSTATUS_TOKEN= \ - ghcr.io/openstatushq/probe:latest + ghcr.io/openstatushq/private-location:latest ``` Replace `` with the token from Step 1. @@ -52,7 +52,7 @@ You should see the container listed with a status of `Up`: ``` CONTAINER ID IMAGE STATUS NAMES -abc123 ghcr.io/openstatushq/probe:latest Up 2 minutes openstatus-probe +abc123 ghcr.io/openstatushq/private-location:latest Up 2 minutes openstatus-probe ``` ## Step 3: assign monitors to your private location @@ -105,7 +105,7 @@ docker run -d \ --restart unless-stopped \ --network host \ -e OPENSTATUS_TOKEN= \ - ghcr.io/openstatushq/probe:latest + ghcr.io/openstatushq/private-location:latest ``` ## What's next -- 2.51.2