MCPcopy Create free account
hub / github.com/PostHog/duckgres / NewOrgRouter

Function NewOrgRouter

controlplane/org_router.go:49–96  ·  view source on GitHub ↗

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))

Source from the content-addressed store, hash-verified

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.
55func 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)

Callers 1

SetupMultiTenantFunction · 0.85

Calls 6

createOrgStackMethod · 0.95
CreateK8sPoolFunction · 0.85
AddMethod · 0.80
SnapshotMethod · 0.80
HealthCheckLoopMethod · 0.65
ErrorMethod · 0.45

Tested by

no test coverage detected