NewOrgRouter creates an OrgRouter from the initial config snapshot.
(store *configstore.ConfigStore, baseCfg K8sWorkerPoolConfig, globalCfg ControlPlaneConfig, srv *server.Server, stsBroker *STSBroker, userSecrets *CPUserSecretManager, resolveDucklingStatus func(context.Context, string) (*provisioner.DucklingStatus, error))
| 47 | |
| 48 | // migrating tracks which orgs have a DuckLake migration in progress. |
| 49 | // During migration, new connections for the org are rejected with a |
| 50 | // retry-friendly error instead of timing out waiting for a worker. |
| 51 | migrating sync.Map // orgID (string) → struct{} |
| 52 | } |
| 53 | |
| 54 | // NewOrgRouter creates an OrgRouter from the initial config snapshot. |
| 55 | func NewOrgRouter(store *configstore.ConfigStore, baseCfg K8sWorkerPoolConfig, globalCfg ControlPlaneConfig, srv *server.Server, stsBroker *STSBroker, userSecrets *CPUserSecretManager, resolveDucklingStatus func(context.Context, string) (*provisioner.DucklingStatus, error)) (*OrgRouter, error) { |
| 56 | tr := &OrgRouter{ |
| 57 | orgs: make(map[string]*OrgStack), |
| 58 | configStore: store, |
| 59 | baseCfg: baseCfg, |
| 60 | globalCfg: globalCfg, |
| 61 | srv: srv, |
| 62 | stsBroker: stsBroker, |
| 63 | userSecrets: userSecrets, |
| 64 | resolveDucklingStatus: resolveDucklingStatus, |
| 65 | } |
| 66 | |
| 67 | sharedCfg := baseCfg |
| 68 | sharedCfg.OrgID = "" |
| 69 | sharedCfg.WorkerIDGenerator = func() int { |
| 70 | return int(tr.nextWorkerID.Add(1)) |
| 71 | } |
| 72 | sharedCfg.RuntimeStore = store |
| 73 | |
| 74 | sharedPoolIface, err := CreateK8sPool(sharedCfg) |
| 75 | if err != nil { |
| 76 | return nil, err |
| 77 | } |
| 78 | sharedPool, ok := sharedPoolIface.(*K8sWorkerPool) |
| 79 | if !ok { |
| 80 | return nil, fmt.Errorf("expected shared K8s pool, got %T", sharedPoolIface) |
| 81 | } |
| 82 | tr.sharedPool = sharedPool |
| 83 | tr.admissionReclaimer = NewAdmissionReclaimer(store, AdmissionReclaimerConfig{ |
| 84 | MaxReservations: globalCfg.AdmissionReclaimerMaxReservations, |
| 85 | }) |
| 86 | |
| 87 | sharedCtx, sharedCancel := context.WithCancel(context.Background()) |
| 88 | tr.sharedCancel = sharedCancel |
| 89 | go tr.sharedPool.HealthCheckLoop(sharedCtx, tr.globalCfg.HealthCheckInterval, tr.onSharedWorkerCrash, tr.onSharedWorkerProgress) |
| 90 | |
| 91 | snap := store.Snapshot() |
| 92 | for _, tc := range snap.Orgs { |
| 93 | // Only create stacks for orgs with ready warehouses (or no warehouse at all for backwards compat) |
| 94 | if tc.Warehouse != nil && tc.Warehouse.State != configstore.ManagedWarehouseStateReady { |
| 95 | slog.Info("Skipping org stack creation (warehouse not ready).", "org", tc.Name, "state", tc.Warehouse.State) |
| 96 | continue |
| 97 | } |
| 98 | if _, err := tr.createOrgStack(tc); err != nil { |
| 99 | slog.Error("Failed to create org stack.", "org", tc.Name, "error", err) |
no test coverage detected