Repository navigation
Add scale down to HistoryPubnetParallelCatchupV2 #433
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
473acdb
9285bf3
d9d9eb6
233285a
6b4f927
ed59cf0
cfe2a5e
e29f2f2
56c99f5
520696c
960cb34
753f255
8e9bcf2
e178fc2
95851b0
c489d09
040e465
8a186ab
72a5baa
e2ef41f
9de0692
e79e2da
8046411
775e515
9ae0847
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -54,38 +54,39 @@ let private sweepKind | |
|
|
||
| // Delete resources older than `cutoff` in the given namespace. Targets the same | ||
| // resource type set the retired `clean` verb did (Service, ConfigMap, | ||
| // StatefulSet, Ingress, Job, DaemonSet, Deployment) and also helm-uninstalls | ||
| // any `parallel-catchup-*` releases (PCv2) so helm's release secrets get | ||
| // tidied along with the workloads. | ||
| // StatefulSet, Ingress, Job, DaemonSet, Deployment) plus ownerless Pods, and | ||
| // also helm-uninstalls any `parallel-catchup-*` releases (PCv2) so helm's | ||
| // release secrets get tidied along with the workloads. | ||
| // | ||
| // Used by both the automatic on-startup orphan sweep (cutoff = now - 2 days) | ||
| // and the explicit `force-clean-namespace` verb (cutoff = DateTime.MaxValue, | ||
| // i.e. everything). | ||
| let private sweepWithCutoff (cutoff: DateTime) (kube: Kubernetes) (ns: string) (apiRateLimit: int) : unit = | ||
| LogInfo "Orphan sweep starting: namespace=%s cutoff=%s" ns (cutoff.ToString("o")) | ||
|
|
||
| // 1. helm-uninstall PCv2 releases whose StatefulSet is older than the cutoff. | ||
| // 1. helm-uninstall PCv2 releases whose worker pod is older than the cutoff. | ||
| // Done before the kubectl-style deletes so helm's release secrets get | ||
| // cleaned up properly, rather than left dangling pointing at deleted | ||
| // resources. | ||
| ApiRateLimit.sleepUntilNextRateLimitedApiCallTime apiRateLimit | ||
|
|
||
| let stsItems = kube.ListNamespacedStatefulSet(namespaceParameter = ns).Items | ||
|
|
||
| for sts in stsItems do | ||
| if isOlderThan cutoff sts.Metadata then | ||
| let name = sts.Metadata.Name | ||
|
|
||
| if name.StartsWith("parallel-catchup-") && name.EndsWith("-stellar-core") then | ||
| let release = name.Substring(0, name.Length - "-stellar-core".Length) | ||
| LogInfo "Orphan sweep: helm uninstall %s" release | ||
|
|
||
| RunShellCommand [| "helm" | ||
| "uninstall" | ||
| release | ||
| "-n" | ||
| ns |] | ||
| |> ignore | ||
| let podItems = kube.ListNamespacedPod(namespaceParameter = ns).Items | ||
|
|
||
|
|
||
| for release in podItems | ||
| |> Seq.filter (fun pod -> isOlderThan cutoff pod.Metadata) | ||
| |> Seq.map (fun pod -> pod.Metadata.Name) | ||
| |> Seq.filter (fun name -> name.StartsWith("parallel-catchup-") && name.Contains("-stellar-core")) | ||
| |> Seq.map (fun name -> name.Substring(0, name.LastIndexOf("-stellar-core"))) | ||
| |> Set.ofSeq do | ||
| LogInfo "Orphan sweep: helm uninstall %s" release | ||
|
|
||
| RunShellCommand [| "helm" | ||
| "uninstall" | ||
| release | ||
| "-n" | ||
| ns |] | ||
| |> ignore | ||
|
|
||
| // 2. Delete the same resource type set the retired `clean` verb targeted. | ||
| // The order matches the old NamespaceContent.Cleanup so dependent | ||
|
|
@@ -111,6 +112,19 @@ let private sweepWithCutoff (cutoff: DateTime) (kube: Kubernetes) (ns: string) ( | |
| |> ignore) | ||
| cutoff | ||
|
|
||
| // PCv2 workers are bare pods, so nothing else in this sweep would reap them. | ||
| // Restricted to ownerless pods: everything else goes away with its controller. | ||
| // Listed here, after the helm uninstalls above, so pods they already removed are not deleted again. | ||
| sweepKind | ||
|
Jonathan-Eid marked this conversation as resolved.
|
||
| apiRateLimit | ||
| "Pod" | ||
| (fun () -> | ||
| kube.ListNamespacedPod(namespaceParameter = ns).Items | ||
| |> Seq.filter (fun p -> isNull p.Metadata.OwnerReferences || p.Metadata.OwnerReferences.Count = 0) | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. This deletes every ownerless pod in the namespace that is older than 2 days, not just PCv2 workers. The sweep runs at every mission start, so any Generated by Claude Code
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. The goal of the sweep is to keep the supercluster namespace completely clean, so far we've had no need to keep recurring resources for a next run, so the extra crossfire would be okay. Also, these changes would be the first to introduce ownerless pods in this namespace, as no other mission spins up bare pods as spinning up bare pods without a deployment or statefulset, etc is unconventional. |
||
| |> Seq.map (fun p -> p.Metadata)) | ||
| (fun n -> kube.DeleteNamespacedPod(namespaceParameter = ns, name = n) |> ignore) | ||
| cutoff | ||
|
|
||
| sweepKind | ||
| apiRateLimit | ||
| "Deployment" | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -1,5 +1,6 @@ | ||
| import os | ||
| import redis | ||
| import socket | ||
| import requests | ||
| import json | ||
| import sys | ||
|
|
@@ -26,6 +27,8 @@ | |
| PROGRESS_QUEUE = os.getenv('PROGRESS_QUEUE', 'in_progress') #LIST | ||
| METRICS = os.getenv('METRICS', 'metrics') # SET | ||
| JOB_OWNERS = os.getenv('JOB_OWNERS', 'job_owners') # HASH | ||
| RETIRING = os.getenv('RETIRING', 'retiring') # SET | ||
| MIN_UNMARKED_WORKERS = int(os.getenv('MIN_UNMARKED_WORKERS', 8)) | ||
| WORKER_PREFIX = os.getenv('WORKER_PREFIX', 'stellar-core') | ||
| NAMESPACE = os.getenv('NAMESPACE', 'default') | ||
| WORKER_COUNT = int(os.getenv('WORKER_COUNT', 3)) | ||
|
|
@@ -77,6 +80,7 @@ def get_logging_level(): | |
| 'jobs_failed': [], | ||
| 'jobs_in_progress': [], | ||
| 'workers': [], | ||
| 'retirable': [], | ||
| 'workers_up': 0, | ||
| 'workers_down': 0, | ||
| 'workers_refresh_duration': 0, | ||
|
|
@@ -148,6 +152,14 @@ def ping_worker(pod_name, retries=1): | |
| time.sleep(STUCK_JOB_PING_DELAY_SECS) | ||
| return False | ||
|
|
||
| def pod_exists(pod_name): | ||
| # Trailing dot skips the resolv.conf search list, so a miss costs one query. | ||
| try: | ||
| socket.gethostbyname(f"{pod_name}.{WORKER_PREFIX}.{NAMESPACE}.svc.cluster.local.") | ||
| return True | ||
| except socket.gaierror: | ||
| return False | ||
|
|
||
| def update_status_and_metrics(): | ||
| global status | ||
| mission_start_time = time.time() | ||
|
|
@@ -174,8 +186,27 @@ def update_status_and_metrics(): | |
| queue_in_progress_count = len(jobs_in_progress) | ||
| queue_remain_count = redis_client.llen(JOB_QUEUE) | ||
|
|
||
| # --- Phase 1c: Mark surplus idle workers as retiring; worker.sh stops claiming once marked --- | ||
| # Names here are pod names ("{WORKER_PREFIX}-{i}") as worker.sh stores them in JOB_OWNERS. | ||
| busy = set(job_owners.values()) | ||
| retiring = redis_client.smembers(RETIRING) | ||
| candidates = sorted(w for w in (f"{WORKER_PREFIX}-{i}" for i in range(WORKER_COUNT)) | ||
| if w not in busy and w not in retiring) | ||
| outstanding = queue_remain_count + queue_in_progress_count | ||
| keep = max(outstanding, MIN_UNMARKED_WORKERS) | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. If the unmarked reserve (as few as Generated by Claude Code
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Yes, this is why the default for MIN_UNMARKED_WORKERS left some headroom to account for some nodes going down, but it's not expected in an ondemand setup for nodes to go down and be reclaimed. That's what led me to the tradeoff of using n bare pods instead of n statefulsets so less load can be applied on the k8s api server. |
||
| to_mark = [] | ||
| # Resolve only when the name count says something could be marked: at t=0 outstanding >= fleet, so zero lookups. | ||
| if len(busy - retiring) + len(candidates) > keep: | ||
| unmarked_idle = [w for w in candidates if pod_exists(w)] | ||
| to_mark = unmarked_idle[:max(0, len(busy - retiring) + len(unmarked_idle) - keep)] | ||
| if to_mark: | ||
| redis_client.sadd(RETIRING, *to_mark) | ||
| logger.info("Marked %d workers retiring (%d outstanding)", len(to_mark), outstanding) | ||
| # Marked on an earlier pass and still idle; names stay here after the driver deletes them. | ||
| retirable = sorted(retiring - busy) | ||
|
|
||
| # --- Phase 2: Quick single-ping check of workers that own in-progress jobs --- | ||
| active_workers = set(job_owners.values()) | ||
| active_workers = busy | ||
| worker_statuses = [] | ||
| workers_up = 0 | ||
| workers_down = 0 | ||
|
|
@@ -231,6 +262,7 @@ def update_status_and_metrics(): | |
| 'jobs_failed': jobs_failed, | ||
| 'jobs_in_progress': jobs_in_progress, | ||
| 'workers': worker_statuses, | ||
| 'retirable': retirable, | ||
| 'workers_up': workers_up, | ||
| 'workers_down': workers_down, | ||
| 'workers_refresh_duration': workers_refresh_duration, | ||
|
|
||
Uh oh!
There was an error while loading. Please reload this page.