diff --git a/p2p/kademlia/network.go b/p2p/kademlia/network.go index 366f0f4a..a2dc1e45 100644 --- a/p2p/kademlia/network.go +++ b/p2p/kademlia/network.go @@ -32,14 +32,22 @@ const ( defaultMaxPayloadSize = 200 // MB errorBusy = "Busy" maxConcurrentFindBatchValsRequests = 25 - defaultExecTimeout = 10 * time.Second + // defaultExecTimeout is for small control-plane RPCs; large payload RPCs + // have explicit entries in execTimeouts below. + defaultExecTimeout = 10 * time.Second ) // Global map for message type timeouts var execTimeouts map[int]time.Duration func init() { - // Initialize the request execution timeout values + // Initialize the request execution timeout values. + // These defaults are intentionally conservative to accommodate slower + // peers and larger payloads. If future deployments consistently see + // responsive nodes, consider reducing the larger RPC timeouts (e.g., + // BatchStoreData/BatchGetValues to ~45s, BatchFindValues to ~30s) to + // fail fast on degraded nodes. Long, user-dependent operations like + // uploads/downloads are governed at higher layers. execTimeouts = map[int]time.Duration{ Ping: 5 * time.Second, FindNode: 10 * time.Second, diff --git a/pkg/common/task/worker.go b/pkg/common/task/worker.go index 280b5fb8..a62752d9 100644 --- a/pkg/common/task/worker.go +++ b/pkg/common/task/worker.go @@ -109,10 +109,15 @@ func NewWorker() *Worker { } // cleanupLoop periodically removes tasks that are in a final state for a grace period +// or any task that has been around for too long func (worker *Worker) cleanupLoop(ctx context.Context) { const ( cleanupInterval = 30 * time.Second finalTaskTTL = 2 * time.Minute + // maxTaskAge removes any task entry after this age, regardless of state. + // Keep greater than the largest server-side task envelope (RegisterTimeout ~75m) + // to avoid pruning legitimate long-running tasks from the worker registry. + maxTaskAge = 2 * time.Hour ) ticker := time.NewTicker(cleanupInterval) @@ -129,11 +134,17 @@ func (worker *Worker) cleanupLoop(ctx context.Context) { kept := worker.tasks[:0] for _, t := range worker.tasks { st := t.Status() - if st != nil && st.SubStatus != nil && st.SubStatus.IsFinal() { - if now.Sub(st.CreatedAt) >= finalTaskTTL { - // drop this finalized task + if st != nil { + // Remove any task older than 30 minutes, regardless of state + if now.Sub(st.CreatedAt) >= maxTaskAge { continue } + // Also remove final tasks after 2 minutes + if st.SubStatus != nil && st.SubStatus.IsFinal() { + if now.Sub(st.CreatedAt) >= finalTaskTTL { + continue + } + } } kept = append(kept, t) } diff --git a/sdk/adapters/supernodeservice/adapter.go b/sdk/adapters/supernodeservice/adapter.go index e55cc845..8008b277 100644 --- a/sdk/adapters/supernodeservice/adapter.go +++ b/sdk/adapters/supernodeservice/adapter.go @@ -370,10 +370,14 @@ func (a *cascadeAdapter) CascadeSupernodeDownload( opts ...grpc.CallOption, ) (*CascadeSupernodeDownloadResponse, error) { - // Use provided context as-is (no correlation IDs) + // Use provided context as-is (no correlation IDs). Add watchdogs: + // - idle timer: reset on every received message (event or chunk). + // - max timer: hard cap for one attempt. + phaseCtx, phaseCancel := context.WithCancel(ctx) + defer phaseCancel() // 1. Open gRPC stream (server-stream) - stream, err := a.client.Download(ctx, &cascade.DownloadRequest{ + stream, err := a.client.Download(phaseCtx, &cascade.DownloadRequest{ ActionId: in.ActionID, Signature: in.Signature, }, opts...) @@ -401,6 +405,23 @@ func (a *cascadeAdapter) CascadeSupernodeDownload( chunkIndex int ) + // 3. Receive streamed responses with liveness watchdog + // Start with a generous prep idle timeout; tighten after first message + currentIdle := downloadPrepIdleTimeout + idleTimer := time.AfterFunc(currentIdle, func() { + a.logger.Error(ctx, "download idle timeout; cancelling stream", "action_id", in.ActionID) + phaseCancel() + }) + defer idleTimer.Stop() + maxTimer := time.AfterFunc(downloadMaxTimeout, func() { + a.logger.Error(ctx, "download max timeout; cancelling stream", "action_id", in.ActionID) + phaseCancel() + }) + defer maxTimer.Stop() + start := time.Now() + lastActivity := start + firstMsg := false + // 3. Receive streamed responses for { resp, err := stream.Recv() @@ -408,6 +429,17 @@ func (a *cascadeAdapter) CascadeSupernodeDownload( break } if err != nil { + // Classify timeouts for clearer upstream handling + if phaseCtx.Err() != nil { + sinceLast := time.Since(lastActivity) + sinceStart := time.Since(start) + switch { + case sinceLast >= downloadIdleTimeout: + return nil, fmt.Errorf("download idle timeout: %w", context.DeadlineExceeded) + case sinceStart >= downloadMaxTimeout: + return nil, fmt.Errorf("download overall timeout: %w", context.DeadlineExceeded) + } + } return nil, fmt.Errorf("stream recv: %w", err) } @@ -415,6 +447,15 @@ func (a *cascadeAdapter) CascadeSupernodeDownload( // 3a. Progress / event message case *cascade.DownloadResponse_Event: + // On first message, tighten idle window for active transfer + if !firstMsg { + firstMsg = true + currentIdle = downloadIdleTimeout + } + if idleTimer != nil { + idleTimer.Reset(currentIdle) + } + lastActivity = time.Now() a.logger.Info(ctx, "supernode event", "event_type", x.Event.EventType, "message", x.Event.Message, "action_id", in.ActionID) if in.EventLogger != nil { @@ -426,10 +467,19 @@ func (a *cascadeAdapter) CascadeSupernodeDownload( }) } - // 3b. Actual data chunk + // 3b. Actual data chunk case *cascade.DownloadResponse_Chunk: data := x.Chunk.Data if len(data) == 0 { + // Treat empty chunks as keep-alive; reset idle + if !firstMsg { + firstMsg = true + currentIdle = downloadIdleTimeout + } + if idleTimer != nil { + idleTimer.Reset(currentIdle) + } + lastActivity = time.Now() continue } if _, err := outFile.Write(data); err != nil { @@ -438,7 +488,14 @@ func (a *cascadeAdapter) CascadeSupernodeDownload( bytesWritten += int64(len(data)) chunkIndex++ - + if !firstMsg { + firstMsg = true + currentIdle = downloadIdleTimeout + } + if idleTimer != nil { + idleTimer.Reset(currentIdle) + } + lastActivity = time.Now() a.logger.Debug(ctx, "received chunk", "chunk_index", chunkIndex, "chunk_size", len(data), "bytes_written", bytesWritten) } } diff --git a/sdk/adapters/supernodeservice/timeouts.go b/sdk/adapters/supernodeservice/timeouts.go index 83311185..884ca83b 100644 --- a/sdk/adapters/supernodeservice/timeouts.go +++ b/sdk/adapters/supernodeservice/timeouts.go @@ -9,3 +9,16 @@ const cascadeUploadTimeout = 60 * time.Minute // cascadeProcessingTimeout bounds the time waiting for server-side processing // and final response (e.g., tx hash) after upload completes. const cascadeProcessingTimeout = 10 * time.Minute + +// Download timeouts (adapter-level) +// - downloadPrepIdleTimeout: idle window before the first message arrives, +// allowing the server time to prepare (e.g., reconstruct large files). +// - downloadIdleTimeout: cancels if no messages (events/chunks) are received +// after transfer begins; protects against stalls while allowing long transfers. +// - downloadMaxTimeout: hard cap for a single download attempt. +const ( + // Give server prep up to ~5m + cushion without client cancelling. + downloadPrepIdleTimeout = 6 * time.Minute + downloadIdleTimeout = 2 * time.Minute + downloadMaxTimeout = 60 * time.Minute +) diff --git a/sdk/docs/cascade-timeouts.md b/sdk/docs/cascade-timeouts.md index b84348c1..8a2af30e 100644 --- a/sdk/docs/cascade-timeouts.md +++ b/sdk/docs/cascade-timeouts.md @@ -1,134 +1,60 @@ -# Cascade Registration Timeouts and Networking - -This document explains how timeouts and deadlines are applied across the SDK cascade registration flow, including the current split between upload and processing phases and the relevant client/server defaults. - -## Purpose - -- Make slow, user‑network–dependent uploads more tolerant without impacting other stages. -- Keep health checks and connection establishment responsive. -- Enable clearer error categorization: upload vs processing. - -## TL;DR Defaults - -- Upload timeout (adapter): `cascadeUploadTimeout = 60m` — covers client-side file streaming to the supernode. -- Processing timeout (adapter): `cascadeProcessingTimeout = 10m` — covers waiting for server progress/final tx hash after upload completes. -- Health check to supernodes (task): `connectionTimeout = 10s` — per-node probe during discovery. -- gRPC connect (client): - - Adds a default `30s` deadline if caller context has none. - - Connection readiness gate: `ConnWaitTime = 10s` per attempt, with `MaxRetries = 3` and retry backoff. -- ALTS handshake (secure transport): `30s` internal read timeouts (client and server sides). -- Supernode gRPC server: - - No per‑RPC timeout for `Register`/`Download` handlers. - - Keepalive is permissive (idle ping at 1h, ping ack timeout 30m). - - Stream tuning: 16MB message caps, 16MB stream window, 160MB conn window, ~20 concurrent streams. +# Cascade Timeouts — Quick Guide -## Control Flow and Contexts +Concise overview of timeout locations, defaults, and intent. -1) `sdk/action/client.go: ClientImpl.StartCascade(ctx, ...)` - - Forwards `ctx` to Task Manager. +## Defaults -2) `sdk/task/manager.go: ManagerImpl.CreateCascadeTask(...)` - - Detaches from caller: `taskCtx := context.WithCancel(context.Background())`. - - All subsequent work uses `taskCtx` (no deadline by default). +- Register (client → server) + - Upload (SDK adapter): 60m — `cascadeUploadTimeout` + - Processing (SDK adapter): 10m — `cascadeProcessingTimeout` + - Server envelope: 75m — `RegisterTimeout` -3) `sdk/task/cascade.go: CascadeTask.Run(ctx)` - - Validates file size; fetches healthy supernodes; registers with one. +- Download (server → client) + - Server preparation: 5m — `DownloadPrepareTimeout` + - Client per‑attempt: 60m — `downloadTimeout` + - Client liveness (SDK adapter): + - Prep idle (pre‑first‑message): 6m — `downloadPrepIdleTimeout` + - Idle (post‑first‑message): 2m — `downloadIdleTimeout` + - Max attempt: 60m — `downloadMaxTimeout` + - Note: File streaming is not server‑bounded; the client governs transfer. -4) Discovery: `sdk/task/task.go: BaseTask.fetchSupernodes` → `BaseTask.isServing` - - `context.WithTimeout(parent, 10s)` for health probe (create client + `HealthCheck`). +- Discovery / Connect + - Health probe per supernode: 10s — `connectionTimeout` + - gRPC connect default: 30s if caller provides no deadline + - Keepalives: permissive (idle ping ~1h, ack timeout ~30m) -5) Registration attempt: `sdk/task/cascade.go: attemptRegistration` - - Client connect: uses task context (no deadline); gRPC injects a 30s default at connect if needed. - - No outer registration timeout here; the adapter handles per‑phase timers. +## Intent and Ordering -6) RPC staging: - - `sdk/net/impl.go: supernodeClient.RegisterCascade` → - - `sdk/adapters/supernodeservice/adapter.go: CascadeSupernodeRegister` performs client‑stream upload and reads server progress / final tx hash. +- Register: server envelope (75m) > SDK phases (60m + 10m) so the client surfaces errors first when appropriate. +- Download: server prep is tight (5m). Transfer is governed by the client with a generous per‑attempt window and two‑phase idle watchdogs. -## Where Timeouts Come From (by Layer) +## Where They Live -- SDK adapter level (registration RPC): - - `cascadeUploadTimeout` (60m): upload phase timer (file chunks + metadata + CloseSend). - - `cascadeProcessingTimeout` (10m): processing phase timer (receive server progress + final tx hash). -- SDK task level: - - `connectionTimeout` (10s): supernode health checks only. +- SDK adapter (upload/download phases): `sdk/adapters/supernodeservice/timeouts.go` +- SDK task (discovery, per‑attempt download): `sdk/task/timeouts.go` +- Supernode service (server envelopes): `supernode/services/cascade/timeouts.go` +- P2P internal RPCs: `p2p/kademlia/network.go` (fixed per‑message timeouts) -- gRPC client (`pkg/net/grpc/client`): - - `defaultTimeout = 30s`: applied to connect if context lacks a deadline. - - `ConnWaitTime = 10s`, `MaxRetries = 3`, backoff configured; keepalives: 30m/30m. +## Notes -- ALTS handshake (`pkg/net/credentials/alts/handshake`): - - `defaultTimeout = 30s` for handshake read operations (client/server). - -- gRPC server (`pkg/net/grpc/server` and supernode runtime): - - No explicit per‑RPC timeouts; generous keepalives; tuned flow control and message sizes for 4MB chunks. - -## SDK Constants - -Timeout constants are defined in dedicated files for clarity: - -- Upload/Processing: `supernode/sdk/adapters/supernodeservice/timeouts.go` -- Connection/health probe: `supernode/sdk/task/timeouts.go` - -Notes: -- `BaseTask.isServing` keeps a short 10s budget for snappy health checks. +- Health checks use a 10s budget for snappy discovery. - gRPC connect/handshake defaults remain unchanged. -## Implementation Details - -The split is implemented inside `CascadeSupernodeRegister` where the phases are naturally separated by the client‑stream CloseSend. - -1) Create a cancelable context from the inbound one for the stream lifetime: - -```go -phaseCtx, cancel := context.WithCancel(ctx) -defer cancel() -stream, err := a.client.Register(phaseCtx, opts...) -``` - -2) Upload phase timer: - -```go -uploadTimer := time.AfterFunc(cascadeUploadTimeout, cancel) - -// send chunks... -// send metadata... - -if err := stream.CloseSend(); err != nil { /* ... */ } -uploadTimer.Stop() -``` - -3) Processing phase timer (server progress → final tx hash): - -```go -processingTimer := time.AfterFunc(cascadeProcessingTimeout, cancel) -defer processingTimer.Stop() - -for { - resp, err := stream.Recv() - // handle EOF, errors, progress, final tx hash -} -``` - -4) Error mapping and events: -- If cancellation occurs during Send loop → classify as upload timeout and emit `SDKUploadTimeout`. -- If cancellation occurs during Recv loop → classify as processing timeout and emit `SDKProcessingTimeout`. -- Surface distinct error messages and publish events accordingly. +## Events (SDK) +- Upload timeout → `SDKUploadFailed` +- Processing timeout → `SDKProcessingTimeout` +- Download failure (timeout/canceled) → `SDKDownloadFailure` This approach requires no request‑struct changes and preserves existing call sites. It uses a single cancelable context across both phases and phase‑specific timers. -## Additional Notes - -- Health checks use `connectionTimeout = 10s` during supernode discovery. -- gRPC client connect behavior: adds a `30s` deadline if none is present, waits up to `ConnWaitTime = 10s` per attempt with retries. -- Downloads use a separate `downloadTimeout = 5m` envelope (per-attempt). On timeout during download, the SDK emits `SDKDownloadFailure` with a reason-coded message `| reason=timeout` and sets `event.KeyMessage = "timeout"`. - -## Operational Guidance +## Minimal Tuning Guidance +- Slow client links: keep download attempt at 60m; adjust idle windows if needed. +- Very large inputs: raise `cascadeUploadTimeout` (keep processing modest at 10m). -- For slow client links: raise `cascadeUploadTimeout` (e.g., 30–120m). Keep processing modest (e.g., 5–10m) unless chain finalization is known to stall. -- Server tuning is already generous; no server change required to support longer uploads. -- Telemetry: differentiate upload vs processing timeout in logs and emitted events for better retry behavior and user messaging. -- Retry policy: on upload timeout, prefer retrying with a different supernode; on processing timeout, consider whether the server might still finalize (idempotency depends on service semantics). +## Reference Map +- SDK: `sdk/task/timeouts.go`, `sdk/adapters/supernodeservice/timeouts.go`, `sdk/adapters/supernodeservice/adapter.go` +- Server: `supernode/services/cascade/timeouts.go`, server handlers in `supernode/node/action/server/cascade` +- Network: `pkg/net/grpc/client`, `p2p/kademlia/network.go` ## File/Code Reference Map @@ -145,7 +71,8 @@ This approach requires no request‑struct changes and preserves existing call s - Supernode - `supernode/supernode/node/supernode/server/server.go` — server options (16MB caps, windows, 20 streams). - - `supernode/supernode/node/action/server/cascade/cascade_action_server.go` — server-side Register/Download handlers (no per‑RPC timeout). + - `supernode/supernode/node/action/server/cascade/cascade_action_server.go` — server-side handlers. + - `supernode/supernode/services/cascade/timeouts.go` — Register (`RegisterTimeout = 75m`) and Download prep (`DownloadPrepareTimeout = 5m`) timeouts. ## Events @@ -160,10 +87,12 @@ This document describes how the SDK applies timeouts and deadlines during cascad - Upload (adapter): `cascadeUploadTimeout = 60m` — client-side streaming of file chunks and metadata. - Processing (adapter): `cascadeProcessingTimeout = 10m` — wait for server progress and final tx hash after upload completes. - Discovery (task): `connectionTimeout = 10s` — per-supernode health probe during discovery. -- Download (task): `downloadTimeout = 5m` — envelope for cascade download. +- Download (task): `downloadTimeout = 60m` — per-attempt envelope. Adapter adds + `downloadPrepIdleTimeout = 6m` (pre-first-message), `downloadIdleTimeout = 2m` + (post-first-message), and `downloadMaxTimeout = 60m`. - gRPC client connect: adds a `30s` deadline if none is present; readiness wait per attempt `ConnWaitTime = 10s` with retries and backoff. - ALTS handshake: internal `30s` read timeouts on both client and server sides. -- Supernode gRPC server: no per-RPC timeout; keepalive is permissive (idle ping ~1h, ack timeout ~30m); flow-control and message-size tuning supports 4MB chunks. +- Supernode gRPC server: task-level timeouts are applied (Register 75m). Download preparation is bounded to 5m; file streaming is client-governed. Keepalive is permissive (idle ping ~1h, ack timeout ~30m); flow-control and message-size tuning supports 4MB chunks. ## Control Flow diff --git a/sdk/task/download.go b/sdk/task/download.go index 95a1fa84..7687503d 100644 --- a/sdk/task/download.go +++ b/sdk/task/download.go @@ -6,7 +6,6 @@ import ( "fmt" "os" "path/filepath" - "time" "github.com/LumeraProtocol/supernode/v2/sdk/adapters/lumera" "github.com/LumeraProtocol/supernode/v2/sdk/adapters/supernodeservice" @@ -14,11 +13,6 @@ import ( "github.com/LumeraProtocol/supernode/v2/sdk/net" ) -// timeouts -const ( - downloadTimeout = 5 * time.Minute -) - type CascadeDownloadTask struct { BaseTask actionId string @@ -270,22 +264,22 @@ func (t *CascadeDownloadTask) attemptConcurrentDownload( // Log failure sn := batch[result.idx] - // Classify failure reason when possible - data := event.EventData{ - event.KeySupernode: sn.GrpcEndpoint, - event.KeySupernodeAddress: sn.CosmosAddress, - event.KeyIteration: baseIteration + result.idx + 1, - event.KeyError: result.err.Error(), - } - msg := "download from super-node failed" - if stderrors.Is(result.err, context.DeadlineExceeded) { - data[event.KeyMessage] = "timeout" - msg += " | reason=timeout" - } else if stderrors.Is(result.err, context.Canceled) { - data[event.KeyMessage] = "canceled" - msg += " | reason=canceled" - } - t.LogEvent(ctx, event.SDKDownloadFailure, msg, data) + // Classify failure reason when possible + data := event.EventData{ + event.KeySupernode: sn.GrpcEndpoint, + event.KeySupernodeAddress: sn.CosmosAddress, + event.KeyIteration: baseIteration + result.idx + 1, + event.KeyError: result.err.Error(), + } + msg := "download from super-node failed" + if stderrors.Is(result.err, context.DeadlineExceeded) { + data[event.KeyMessage] = "timeout" + msg += " | reason=timeout" + } else if stderrors.Is(result.err, context.Canceled) { + data[event.KeyMessage] = "canceled" + msg += " | reason=canceled" + } + t.LogEvent(ctx, event.SDKDownloadFailure, msg, data) errs = append(errs, result.err) case <-ctx.Done(): diff --git a/sdk/task/timeouts.go b/sdk/task/timeouts.go index f6e1e7e6..62ae733d 100644 --- a/sdk/task/timeouts.go +++ b/sdk/task/timeouts.go @@ -2,7 +2,22 @@ package task import "time" -// connectionTimeout bounds supernode health/connection probing. -// Keep this short to preserve snappy discovery without impacting long uploads. -const connectionTimeout = 10 * time.Second +// Connection and health check timeouts +const ( + // connectionTimeout bounds supernode health/connection probing. + // Keep this short to preserve snappy discovery without impacting long uploads. + connectionTimeout = 10 * time.Second +) +// Task execution timeouts +const ( + // downloadTimeout bounds a single download attempt at the task layer. + // This should exceed typical slow-client scenarios; fine-grained + // liveness is enforced by the adapter via an idle timeout. + downloadTimeout = 60 * time.Minute +) + +// Note: Upload and processing timeouts are defined in sdk/adapters/supernodeservice/timeouts.go +// as they are specific to the adapter implementation: +// - cascadeUploadTimeout = 60 * time.Minute (for slow network uploads) +// - cascadeProcessingTimeout = 10 * time.Minute (for server-side processing) diff --git a/sn-manager/cmd/start.go b/sn-manager/cmd/start.go index 0166bc10..de03c6dd 100644 --- a/sn-manager/cmd/start.go +++ b/sn-manager/cmd/start.go @@ -121,6 +121,17 @@ func runStart(cmd *cobra.Command, args []string) error { } } + // Mandatory version sync on startup: ensure both sn-manager and SuperNode + // are at the latest stable release. This bypasses regular updater checks + // (gateway idleness, same-major policy) to guarantee a consistent baseline. + // Runs once before monitoring begins. + func() { + u := updater.New(home, cfg, appVersion) + // Do not block startup on failures; best-effort sync + defer func() { recover() }() + u.ForceSyncToLatest(context.Background()) + }() + // Start auto-updater if enabled var autoUpdater *updater.AutoUpdater if cfg.Updates.AutoUpgrade { diff --git a/sn-manager/internal/updater/updater.go b/sn-manager/internal/updater/updater.go index 02e69ee5..b0e2f1ea 100644 --- a/sn-manager/internal/updater/updater.go +++ b/sn-manager/internal/updater/updater.go @@ -20,7 +20,16 @@ import ( "google.golang.org/protobuf/encoding/protojson" ) -const gatewayTimeout = 15 * time.Second +// Global updater timing constants +const ( + // gatewayTimeout bounds the local gateway status probe + gatewayTimeout = 15 * time.Second + // updateCheckInterval is how often the periodic updater runs + updateCheckInterval = 10 * time.Minute + // forceUpdateAfter is the age threshold after a release is published + // beyond which updates are applied regardless of normal gates (idle, policy) + forceUpdateAfter = 1 * time.Hour +) type AutoUpdater struct { config *config.Config @@ -59,18 +68,17 @@ func (u *AutoUpdater) Start(ctx context.Context) { return } - // Fixed update check interval: 10 minutes - interval := 10 * time.Minute - u.ticker = time.NewTicker(interval) + // Fixed update check interval + u.ticker = time.NewTicker(updateCheckInterval) // Run an immediate check on startup so restarts don't wait a full interval - u.checkAndUpdateCombined() + u.checkAndUpdateCombined(false) go func() { for { select { case <-u.ticker.C: - u.checkAndUpdateCombined() + u.checkAndUpdateCombined(false) case <-u.stopCh: return case <-ctx.Done(): @@ -172,7 +180,17 @@ func (u *AutoUpdater) isGatewayIdle() (bool, bool) { // downloads the release tarball once to update sn-manager and SuperNode. // Order: update sn-manager first (prepare new binary), then SuperNode, then // trigger restart if manager was updated. -func (u *AutoUpdater) checkAndUpdateCombined() { +// ForceSyncToLatest performs a one-shot forced sync to the latest stable +// release, bypassing standard gating checks (gateway idle, same-major policy). +// Intended for mandatory checks at manager start. +func (u *AutoUpdater) ForceSyncToLatest(_ context.Context) { + u.checkAndUpdateCombined(true) +} + +// checkAndUpdateCombined performs a single release check and, if needed, +// downloads the release tarball once to update sn-manager and SuperNode. +// If force is true, bypass gateway idleness and version policy checks. +func (u *AutoUpdater) checkAndUpdateCombined(force bool) { // Fetch latest stable release once release, err := u.githubClient.GetLatestStableRelease() @@ -186,18 +204,34 @@ func (u *AutoUpdater) checkAndUpdateCombined() { return } + // If the latest release has been out for > 4 hours, elevate to force mode + if !force { + if !release.PublishedAt.IsZero() && time.Since(release.PublishedAt) > forceUpdateAfter { + force = true + } + } + // Determine if sn-manager should update (same criteria: stable, same major) managerNeedsUpdate := false ver := strings.TrimSpace(u.managerVersion) if ver != "" && ver != "dev" && !strings.EqualFold(ver, "unknown") { - if utils.SameMajor(ver, latest) && utils.CompareVersions(ver, latest) < 0 { - managerNeedsUpdate = true + if force { + managerNeedsUpdate = !strings.EqualFold(ver, latest) + } else { + if utils.SameMajor(ver, latest) && utils.CompareVersions(ver, latest) < 0 { + managerNeedsUpdate = true + } } } // Determine if SuperNode should update using existing policy currentSN := u.config.Updates.CurrentVersion - supernodeNeedsUpdate := u.ShouldUpdate(currentSN, latest) + supernodeNeedsUpdate := false + if force { + supernodeNeedsUpdate = !strings.EqualFold(strings.TrimPrefix(currentSN, "v"), strings.TrimPrefix(latest, "v")) + } else { + supernodeNeedsUpdate = u.ShouldUpdate(currentSN, latest) + } if !managerNeedsUpdate && !supernodeNeedsUpdate { return @@ -205,14 +239,16 @@ func (u *AutoUpdater) checkAndUpdateCombined() { // Gate all updates (manager + SuperNode) on gateway idleness // to avoid disrupting traffic during a self-update. - if idle, isErr := u.isGatewayIdle(); !idle { - if isErr { - // Track errors and possibly request a clean SuperNode restart - u.handleGatewayError() - } else { - log.Println("Gateway busy, deferring updates") + if !force { + if idle, isErr := u.isGatewayIdle(); !idle { + if isErr { + // Track errors and possibly request a clean SuperNode restart + u.handleGatewayError() + } else { + log.Println("Gateway busy, deferring updates") + } + return } - return } // Download the combined release tarball once diff --git a/supernode/services/cascade/adaptors/p2p.go b/supernode/services/cascade/adaptors/p2p.go index fcaad76a..dc8c077a 100644 --- a/supernode/services/cascade/adaptors/p2p.go +++ b/supernode/services/cascade/adaptors/p2p.go @@ -232,8 +232,15 @@ func (c *p2pImpl) storeSymbolsInP2P(ctx context.Context, taskID, root string, fi return 0, 0, 0, fmt.Errorf("load symbols: %w", err) } - symCtx, cancel := context.WithTimeout(ctx, 5*time.Minute) - defer cancel() + // Timeouts: rely on inner per-RPC limits and outer task envelope + // - Each node RPC invoked by StoreBatch uses Network.Call with per-message + // timeouts (BatchStoreData currently ~75s) enforced in kademlia/network.go. + // - The server-side Register handler is wrapped in a long envelope timeout + // (RegisterTimeout) to guarantee eventual completion/cancellation. + // Therefore, we avoid adding another mid-layer timer here to prevent + // premature cancellation of large batches and to let lower-level retries + // and per-RPC deadlines operate as designed. + symCtx := ctx rate, requests, err := c.p2p.StoreBatch(symCtx, symbols, storage.P2PDataRaptorQSymbol, taskID) if err != nil { diff --git a/supernode/services/cascade/download.go b/supernode/services/cascade/download.go index c403b729..a6851397 100644 --- a/supernode/services/cascade/download.go +++ b/supernode/services/cascade/download.go @@ -1,11 +1,11 @@ package cascade import ( - "bytes" - "context" - "fmt" - "os" - "sort" + "bytes" + "context" + "fmt" + "os" + "sort" actiontypes "github.com/LumeraProtocol/lumera/x/action/v1/types" "github.com/LumeraProtocol/supernode/v2/pkg/codec" @@ -31,11 +31,18 @@ type DownloadResponse struct { DownloadedDir string } +// Download preparation is bounded by DownloadPrepareTimeout (see timeouts.go). +// The subsequent file streaming to the client is not bounded by this timer. + func (task *CascadeRegistrationTask) Download( ctx context.Context, req *DownloadRequest, send func(resp *DownloadResponse) error, ) (err error) { + // Bound the preparation phase only (metadata/layout/symbols/restore) + ctx, cancel := context.WithTimeout(ctx, DownloadPrepareTimeout) + defer cancel() + fields := logtrace.Fields{logtrace.FieldMethod: "Download", logtrace.FieldRequest: req} logtrace.Info(ctx, "cascade-action-download request received", fields) @@ -124,10 +131,10 @@ func (task *CascadeRegistrationTask) downloadArtifacts(ctx context.Context, acti } func (task *CascadeRegistrationTask) restoreFileFromLayout( - ctx context.Context, - layout codec.Layout, - dataHash string, - actionID string, + ctx context.Context, + layout codec.Layout, + dataHash string, + actionID string, ) (string, string, error) { fields := logtrace.Fields{ @@ -139,18 +146,18 @@ func (task *CascadeRegistrationTask) restoreFileFromLayout( } sort.Strings(allSymbols) - totalSymbols := len(allSymbols) - requiredSymbols := (totalSymbols*requiredSymbolPercent + 99) / 100 + totalSymbols := len(allSymbols) + requiredSymbols := (totalSymbols*requiredSymbolPercent + 99) / 100 - fields["totalSymbols"] = totalSymbols - fields["requiredSymbols"] = requiredSymbols - logtrace.Info(ctx, "symbols to be retrieved", fields) + fields["totalSymbols"] = totalSymbols + fields["requiredSymbols"] = requiredSymbols + logtrace.Info(ctx, "symbols to be retrieved", fields) - // Progressive retrieval moved to helper for readability/testing - decodeInfo, err := task.retrieveAndDecodeProgressively(ctx, allSymbols, layout, actionID, fields) - if err != nil { - return "", "", err - } + // Progressive retrieval moved to helper for readability/testing + decodeInfo, err := task.retrieveAndDecodeProgressively(ctx, allSymbols, layout, actionID, fields) + if err != nil { + return "", "", err + } fileHash, err := crypto.HashFileIncrementally(decodeInfo.FilePath, 0) if err != nil { @@ -164,8 +171,8 @@ func (task *CascadeRegistrationTask) restoreFileFromLayout( return "", "", errors.New("file hash is nil") } - // Validate final payload hash against on-chain data hash - err = task.verifyDataHash(ctx, fileHash, dataHash, fields) + // Validate final payload hash against on-chain data hash + err = task.verifyDataHash(ctx, fileHash, dataHash, fields) if err != nil { logtrace.Error(ctx, "failed to verify hash", fields) fields[logtrace.FieldError] = err.Error() diff --git a/supernode/services/cascade/register.go b/supernode/services/cascade/register.go index b9c9de83..640f1a98 100644 --- a/supernode/services/cascade/register.go +++ b/supernode/services/cascade/register.go @@ -24,6 +24,10 @@ type RegisterResponse struct { TxHash string } +// RegisterTimeout bounds the execution time of a Register task to prevent tasks +// from lingering in a non-final state if a dependency stalls. Keep greater than +// the SDK's upload+processing budgets so the client cancels first. + // Register processes the upload request for cascade input data. // 1- Fetch & validate action (it should be a cascade action registered on the chain) // 2- Ensure this super-node is eligible to process the action (should be in the top supernodes list for the action block height) @@ -44,6 +48,9 @@ func (task *CascadeRegistrationTask) Register( req *RegisterRequest, send func(resp *RegisterResponse) error, ) (err error) { + // Defensive envelope deadline to guarantee task finalization + ctx, cancel := context.WithTimeout(ctx, RegisterTimeout) + defer cancel() fields := logtrace.Fields{logtrace.FieldMethod: "Register", logtrace.FieldRequest: req} logtrace.Info(ctx, "cascade-action-registration request received", fields) @@ -145,13 +152,13 @@ func (task *CascadeRegistrationTask) Register( task.streamEvent(SupernodeEventTypeRqIDsVerified, "rq-ids have been verified", "", send) /* 10. Simulate finalize to avoid storing artefacts if it would fail ---------- */ - if _, err := task.LumeraClient.SimulateFinalizeAction(ctx, action.ActionID, rqidResp.RQIDs); err != nil { - fields[logtrace.FieldError] = err.Error() - logtrace.Info(ctx, "finalize action simulation failed", fields) - // Emit explicit simulation failure event for client visibility - task.streamEvent(SupernodeEventTypeFinalizeSimulationFailed, "finalize action simulation failed", "", send) - return task.wrapErr(ctx, "finalize action simulation failed", err, fields) - } + if _, err := task.LumeraClient.SimulateFinalizeAction(ctx, action.ActionID, rqidResp.RQIDs); err != nil { + fields[logtrace.FieldError] = err.Error() + logtrace.Info(ctx, "finalize action simulation failed", fields) + // Emit explicit simulation failure event for client visibility + task.streamEvent(SupernodeEventTypeFinalizeSimulationFailed, "finalize action simulation failed", "", send) + return task.wrapErr(ctx, "finalize action simulation failed", err, fields) + } logtrace.Info(ctx, "finalize action simulation passed", fields) // Transmit as a standard event so SDK can propagate it (dedicated type) task.streamEvent(SupernodeEventTypeFinalizeSimulated, "finalize action simulation passed", "", send) diff --git a/supernode/services/cascade/timeouts.go b/supernode/services/cascade/timeouts.go new file mode 100644 index 00000000..110dba3d --- /dev/null +++ b/supernode/services/cascade/timeouts.go @@ -0,0 +1,21 @@ +package cascade + +import "time" + +// Server-side task envelope timeouts for Cascade service. +// +// Rationale: +// - RegisterTimeout must exceed the total of SDK upload + processing budgets +// so the client surfaces errors; currently upload=60m and processing=10m. +// - Download uses a split approach: a tight server-side preparation timeout +// and a relaxed client-governed transfer window. +const ( + // RegisterTimeout bounds the entire Register RPC handler lifetime. + // Must be greater than SDK's upload (60m) + processing (10m) budgets. + RegisterTimeout = 75 * time.Minute + + // DownloadPrepareTimeout bounds the server-side preparation phase for + // downloads (fetch metadata, retrieve symbols, reconstruct and verify file). + // This phase is independent of client bandwidth and should be quick. + DownloadPrepareTimeout = 5 * time.Minute +) diff --git a/supernode/services/common/base/supernode_task.go b/supernode/services/common/base/supernode_task.go index 937e6013..6cde3d4f 100644 --- a/supernode/services/common/base/supernode_task.go +++ b/supernode/services/common/base/supernode_task.go @@ -25,7 +25,21 @@ type SuperNodeTask struct { func (task *SuperNodeTask) RunHelper(ctx context.Context, clean TaskCleanerFunc) error { ctx = task.context(ctx) logtrace.Debug(ctx, "Start task", logtrace.Fields{}) - defer logtrace.Info(ctx, "Task canceled", logtrace.Fields{}) + defer func() { + // Log accurate end-state when task finishes + st := task.Status() + if st != nil && st.SubStatus != nil { + if st.SubStatus.IsFailure() { + logtrace.Info(ctx, "Task canceled", logtrace.Fields{}) + } else if st.SubStatus.IsFinal() { + logtrace.Info(ctx, "Task completed", logtrace.Fields{}) + } else { + logtrace.Info(ctx, "Task ended", logtrace.Fields{}) + } + } else { + logtrace.Info(ctx, "Task ended", logtrace.Fields{}) + } + }() defer task.Cancel() task.SetStatusNotifyFunc(func(status *state.Status) { diff --git a/testnet_version_check.sh b/testnet_version_check.sh index e3a53a4a..e61bf32a 100755 --- a/testnet_version_check.sh +++ b/testnet_version_check.sh @@ -7,7 +7,7 @@ set -o pipefail # # Usage: ./check_versions.sh (no args) -TIMEOUT=10 +TIMEOUT=5 API_URL="https://lcd.testnet.lumera.io/LumeraProtocol/lumera/supernode/list_super_nodes?pagination.limit=1000&pagination.count_total=true" TOTAL_CHECKED=0