Skip to content

Commit d40fd50

Browse files
committed
warpctl: config upstream use consistent routing for stateful services
1 parent 498fece commit d40fd50

1 file changed

Lines changed: 43 additions & 15 deletions

File tree

warpctl/config.go

Lines changed: 43 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -341,6 +341,7 @@ type ServiceConfig struct {
341341
LbExposed *bool `yaml:"lb_exposed,omitempty"`
342342
Websocket *bool `yaml:"websocket,omitempty"`
343343
Streamable *bool `yaml:"streamable,omitempty"`
344+
Stateful *bool `yaml:"stateful,omitempty"`
344345
Hosts []string `yaml:"hosts,omitempty"`
345346
EnvVars map[string]string `yaml:"env_vars,omitempty"`
346347
Mount map[string]string `yaml:"mount,omitempty"`
@@ -399,6 +400,11 @@ func (self *ServiceConfig) isStreamable() bool {
399400
return self.Streamable != nil && *self.Streamable
400401
}
401402

403+
func (self *ServiceConfig) isStateful() bool {
404+
// default false
405+
return self.Stateful != nil && *self.Stateful
406+
}
407+
402408
func (self *ServiceConfig) memoryLimit() (memoryLimit ByteCount) {
403409
if self.MemoryLimit == "" {
404410
return
@@ -1799,8 +1805,11 @@ func (self *NginxConfig) addUpstreamBlocks() {
17991805
// service-block-<service>
18001806
// service-block-<service>-<block>
18011807
for _, service := range self.services() {
1808+
servicesConfig := self.servicesConfig.Versions[0]
1809+
serviceConfig := servicesConfig.Services[service]
1810+
18021811
// only service port 80 is exposed via the html block
1803-
if !slices.Contains(self.servicesConfig.Versions[0].Services[service].TcpPorts(), 80) {
1812+
if !slices.Contains(serviceConfig.TcpPorts(), 80) {
18041813
continue
18051814
}
18061815

@@ -1811,7 +1820,7 @@ func (self *NginxConfig) addUpstreamBlocks() {
18111820
continue
18121821
}
18131822

1814-
keepalive := self.servicesConfig.Versions[0].Services[service].Keepalive
1823+
keepalive := serviceConfig.Keepalive
18151824
if keepalive == nil {
18161825
keepalive = self.lbBlockInfo.lbBlock.Keepalive
18171826
if keepalive == nil {
@@ -1827,9 +1836,15 @@ func (self *NginxConfig) addUpstreamBlocks() {
18271836
)
18281837

18291838
self.block(upstream, func() {
1830-
self.raw(`
1831-
least_conn;
1832-
`)
1839+
if serviceConfig.isStateful() {
1840+
self.raw(`
1841+
hash $remote_addr consistent;
1842+
`)
1843+
} else {
1844+
self.raw(`
1845+
least_conn;
1846+
`)
1847+
}
18331848

18341849
for _, block := range blocks {
18351850
portBlock, ok := httpPortBlocks[service][block][80]
@@ -1840,7 +1855,7 @@ func (self *NginxConfig) addUpstreamBlocks() {
18401855

18411856
upstreamServer := templateString("server {{.dockerNetwork}}:{{.externalPort}} weight={{.weight}} max_fails=0 max_conns=32768;",
18421857
map[string]any{
1843-
"dockerNetwork": self.servicesConfig.Versions[0].ServicesDockerNetwork,
1858+
"dockerNetwork": servicesConfig.ServicesDockerNetwork,
18441859
"externalPort": portBlock.externalPort,
18451860
"weight": blockInfo.weight,
18461861
},
@@ -1878,14 +1893,20 @@ func (self *NginxConfig) addUpstreamBlocks() {
18781893
)
18791894

18801895
self.block(blockUpstream, func() {
1881-
self.raw(`
1882-
least_conn;
1883-
`)
1896+
if serviceConfig.isStateful() {
1897+
self.raw(`
1898+
hash $remote_addr consistent;
1899+
`)
1900+
} else {
1901+
self.raw(`
1902+
least_conn;
1903+
`)
1904+
}
18841905

18851906
self.raw(`
18861907
server {{.dockerNetwork}}:{{.externalPort}};
18871908
`, map[string]any{
1888-
"dockerNetwork": self.servicesConfig.Versions[0].ServicesDockerNetwork,
1909+
"dockerNetwork": servicesConfig.ServicesDockerNetwork,
18891910
"externalPort": portBlock.externalPort,
18901911
})
18911912

@@ -2484,7 +2505,8 @@ func (self *NginxConfig) addStreamUpstreamBlocks() {
24842505
continue
24852506
}
24862507

2487-
serviceConfig := self.servicesConfig.Versions[0].Services[service]
2508+
servicesConfig := self.servicesConfig.Versions[0]
2509+
serviceConfig := servicesConfig.Services[service]
24882510

24892511
for portType, ports := range serviceConfig.AllStreamPorts() {
24902512
for _, port := range ports {
@@ -2499,17 +2521,23 @@ func (self *NginxConfig) addStreamUpstreamBlocks() {
24992521
)
25002522

25012523
self.block(upstream, func() {
2502-
self.raw(`
2503-
least_conn;
2504-
`)
2524+
if serviceConfig.isStateful() {
2525+
self.raw(`
2526+
hash $remote_addr consistent;
2527+
`)
2528+
} else {
2529+
self.raw(`
2530+
least_conn;
2531+
`)
2532+
}
25052533

25062534
for _, block := range blocks {
25072535
for _, portBlock := range streamPortBlocks[service][block] {
25082536
if portBlock.port == port {
25092537
blockInfo := self.blockInfos[service][block]
25102538
upstreamServer := templateString("server {{.dockerNetwork}}:{{.externalPort}} weight={{.weight}} max_fails=0 max_conns=32768;",
25112539
map[string]any{
2512-
"dockerNetwork": self.servicesConfig.Versions[0].ServicesDockerNetwork,
2540+
"dockerNetwork": servicesConfig.ServicesDockerNetwork,
25132541
"externalPort": portBlock.externalPort,
25142542
"weight": blockInfo.weight,
25152543
},

0 commit comments

Comments
 (0)