Skip to content

Commit cd1b345

Browse files
authored
Merge pull request #72 from ConstellationCrypto/inomurko/fix-read-only-fallback
fix: enable read only
2 parents 4d9d54d + 804b2d7 commit cd1b345

8 files changed

Lines changed: 124 additions & 4 deletions

File tree

celestia_server.go

Lines changed: 7 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -99,7 +99,8 @@ type CelestiaServer struct {
9999
metricsServer *http.Server // Separate server for graceful shutdown
100100

101101
// Fallback provider for write-through and read-fallback
102-
fallback fallback.Provider
102+
fallback fallback.Provider
103+
fallbackMode string
103104

104105
// Track pending async operations for graceful shutdown
105106
wg sync.WaitGroup
@@ -119,6 +120,7 @@ func NewCelestiaServer(
119120
metricsEnabled bool,
120121
metricsPort int,
121122
fallbackProvider fallback.Provider,
123+
fallbackMode string,
122124
log log.Logger,
123125
) *CelestiaServer {
124126
endpoint := net.JoinHostPort(host, strconv.Itoa(port))
@@ -138,6 +140,7 @@ func NewCelestiaServer(
138140
metricsEnabled: metricsEnabled,
139141
metricsPort: metricsPort,
140142
fallback: fallbackProvider,
143+
fallbackMode: fallback.NormalizeMode(fallbackMode),
141144
httpServer: &http.Server{
142145
Addr: endpoint,
143146
ReadTimeout: httpReadTimeout,
@@ -278,8 +281,8 @@ func (d *CelestiaServer) getBlob(ctx context.Context, comm []byte) ([]byte, erro
278281

279282
data, celestiaErr := d.store.Get(getCtx, comm)
280283
if celestiaErr == nil {
281-
// Success from Celestia - read-through to fallback for future requests
282-
if d.fallback.Available() {
284+
// Success from Celestia - optionally read-through to fallback for future requests
285+
if d.fallback.Available() && fallback.WriteEnabled(d.fallbackMode) {
283286
d.wg.Add(1)
284287
go d.putFallback(context.Background(), comm, data)
285288
}
@@ -443,7 +446,7 @@ func (d *CelestiaServer) HandlePut(w http.ResponseWriter, r *http.Request) {
443446
"duration", duration)
444447

445448
// Write to fallback provider asynchronously (non-blocking)
446-
if d.fallback.Available() {
449+
if d.fallback.Available() && fallback.WriteEnabled(d.fallbackMode) {
447450
d.wg.Add(1)
448451
go d.putFallback(context.Background(), commitment, input)
449452
}

celestia_server_test.go

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -103,11 +103,16 @@ func createTestServer(t *testing.T, store Store) *CelestiaServer {
103103
false, // metrics disabled for unit tests
104104
0,
105105
nil, // fallback provider (nil = NoopProvider)
106+
fallback.ModeBoth,
106107
logger,
107108
)
108109
}
109110

110111
func createTestServerWithFallback(t *testing.T, store Store, provider fallback.Provider) *CelestiaServer {
112+
return createTestServerWithFallbackMode(t, store, provider, fallback.ModeBoth)
113+
}
114+
115+
func createTestServerWithFallbackMode(t *testing.T, store Store, provider fallback.Provider, mode string) *CelestiaServer {
111116
logger := log.New()
112117

113118
return NewCelestiaServer(
@@ -123,6 +128,7 @@ func createTestServerWithFallback(t *testing.T, store Store, provider fallback.P
123128
false,
124129
0,
125130
provider,
131+
mode,
126132
logger,
127133
)
128134
}
@@ -338,6 +344,23 @@ func TestGetBlob_DoesNotUseFallbackOnCelestiaNotFound(t *testing.T) {
338344
assert.Equal(t, 0, fallbackProvider.getCalls)
339345
}
340346

347+
func TestGetBlob_ReadFallbackMode_SkipsWriteThrough(t *testing.T) {
348+
blob := []byte("blob from celestia")
349+
fallbackProvider := &mockFallbackProvider{available: true}
350+
store := &mockStore{
351+
getFunc: func(ctx context.Context, key []byte) ([]byte, error) {
352+
return blob, nil
353+
},
354+
}
355+
356+
server := createTestServerWithFallbackMode(t, store, fallbackProvider, fallback.ModeReadFallback)
357+
358+
data, err := server.getBlob(context.Background(), commitmentForData(t, blob))
359+
require.NoError(t, err)
360+
assert.Equal(t, blob, data)
361+
assert.Equal(t, 0, fallbackProvider.putCalls)
362+
}
363+
341364
func TestGetBlob_VerifiesFallbackDataAgainstCommitment(t *testing.T) {
342365
commitment := commitmentForData(t, []byte("expected blob"))
343366
store := &mockStore{

cmd/daserver/config.go

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@ import (
99
"time"
1010

1111
celestia "github.com/celestiaorg/op-alt-da"
12+
"github.com/celestiaorg/op-alt-da/fallback"
1213
"github.com/celestiaorg/op-alt-da/signer"
1314
)
1415

@@ -119,6 +120,9 @@ type FallbackConfig struct {
119120
// Provider type: "s3"
120121
Provider string `toml:"provider"`
121122

123+
// Mode controls fallback write/read behavior: write_through, read_fallback, or both.
124+
Mode string `toml:"mode"`
125+
122126
// S3 configuration (when provider = "s3")
123127
S3 S3Config `toml:"s3"`
124128
}
@@ -183,6 +187,7 @@ func DefaultConfig() Config {
183187
Fallback: FallbackConfig{
184188
Enabled: false,
185189
Provider: "s3",
190+
Mode: fallback.ModeBoth,
186191
S3: S3Config{
187192
Region: "us-east-1",
188193
Timeout: "30s",
@@ -252,6 +257,10 @@ func (c *Config) Validate() error {
252257
}
253258
}
254259

260+
if err := fallback.ValidateMode(c.Fallback.Mode); err != nil {
261+
return err
262+
}
263+
255264
return nil
256265
}
257266

cmd/daserver/config_builder.go

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -150,6 +150,9 @@ func BuildConfigFromCLI(ctx *cli.Context) (*Config, error) {
150150
if ctx.IsSet(FallbackProviderFlagName) {
151151
cfg.Fallback.Provider = ctx.String(FallbackProviderFlagName)
152152
}
153+
if ctx.IsSet(FallbackModeFlagName) {
154+
cfg.Fallback.Mode = ctx.String(FallbackModeFlagName)
155+
}
153156
if ctx.IsSet(FallbackS3BucketFlagName) {
154157
cfg.Fallback.S3.Bucket = ctx.String(FallbackS3BucketFlagName)
155158
}

cmd/daserver/entrypoint.go

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -121,6 +121,7 @@ func StartDAServer(cliCtx *cli.Context) error {
121121

122122
// Initialize fallback provider (prefer TOML config, fall back to CLI flags)
123123
var fallbackProvider fallback.Provider = &fallback.NoopProvider{}
124+
fallbackMode := firstNonEmpty(cfg.Fallback.Mode, cliCtx.String(FallbackModeFlagName))
124125
fallbackEnabled := cfg.Fallback.Enabled || cliCtx.Bool(FallbackEnabledFlagName)
125126
if fallbackEnabled {
126127
provider := cfg.Fallback.Provider
@@ -162,10 +163,12 @@ func StartDAServer(cliCtx *cli.Context) error {
162163
return fmt.Errorf("failed to initialize S3 fallback provider: %w", err)
163164
}
164165
fallbackProvider = s3Provider
166+
fallbackMode := firstNonEmpty(cfg.Fallback.Mode, cliCtx.String(FallbackModeFlagName))
165167
l.Info("Fallback provider initialized",
166168
"provider", "s3",
167169
"bucket", s3Cfg.Bucket,
168170
"prefix", s3Cfg.Prefix,
171+
"mode", fallback.NormalizeMode(fallbackMode),
169172
"read_legacy_blobs", s3Cfg.ReadLegacyBlobs)
170173
default:
171174
return fmt.Errorf("unknown fallback provider: %s", provider)
@@ -185,6 +188,7 @@ func StartDAServer(cliCtx *cli.Context) error {
185188
cfg.Metrics.Enabled,
186189
cfg.Metrics.Port,
187190
fallbackProvider,
191+
fallbackMode,
188192
l,
189193
)
190194
} else {

cmd/daserver/flags.go

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ import (
1010
"github.com/urfave/cli/v2"
1111

1212
celestia "github.com/celestiaorg/op-alt-da"
13+
"github.com/celestiaorg/op-alt-da/fallback"
1314
"github.com/celestiaorg/op-alt-da/fallback/s3"
1415
"github.com/celestiaorg/op-alt-da/signer"
1516
opservice "github.com/ethereum-optimism/optimism/op-service"
@@ -47,6 +48,7 @@ const (
4748
// fallback provider flags
4849
FallbackEnabledFlagName = "fallback.enabled"
4950
FallbackProviderFlagName = "fallback.provider"
51+
FallbackModeFlagName = "fallback.mode"
5052
FallbackS3BucketFlagName = "fallback.s3.bucket"
5153
FallbackS3PrefixFlagName = "fallback.s3.prefix"
5254
FallbackS3EndpointFlagName = "fallback.s3.endpoint"
@@ -193,6 +195,12 @@ var (
193195
Value: "s3",
194196
EnvVars: prefixEnvVars("FALLBACK_PROVIDER"),
195197
}
198+
FallbackModeFlag = &cli.StringFlag{
199+
Name: FallbackModeFlagName,
200+
Usage: "Fallback mode: write_through, read_fallback, or both",
201+
Value: fallback.ModeBoth,
202+
EnvVars: prefixEnvVars("FALLBACK_MODE"),
203+
}
196204
FallbackS3BucketFlag = &cli.StringFlag{
197205
Name: FallbackS3BucketFlagName,
198206
Usage: "S3 bucket name for fallback storage",
@@ -277,6 +285,7 @@ var optionalFlags = []cli.Flag{
277285
// Fallback flags
278286
FallbackEnabledFlag,
279287
FallbackProviderFlag,
288+
FallbackModeFlag,
280289
FallbackS3BucketFlag,
281290
FallbackS3PrefixFlag,
282291
FallbackS3EndpointFlag,
@@ -300,6 +309,7 @@ func init() {
300309
type CLIFallbackConfig struct {
301310
Enabled bool
302311
Provider string
312+
Mode string
303313
S3 s3.Config
304314
}
305315

@@ -364,6 +374,7 @@ func ReadCLIConfig(ctx *cli.Context) CLIConfig {
364374
Fallback: CLIFallbackConfig{
365375
Enabled: ctx.Bool(FallbackEnabledFlagName),
366376
Provider: ctx.String(FallbackProviderFlagName),
377+
Mode: ctx.String(FallbackModeFlagName),
367378
S3: s3.Config{
368379
Bucket: ctx.String(FallbackS3BucketFlagName),
369380
Prefix: ctx.String(FallbackS3PrefixFlagName),

fallback/mode.go

Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,42 @@
1+
package fallback
2+
3+
import "fmt"
4+
5+
const (
6+
ModeWriteThrough = "write_through"
7+
ModeReadFallback = "read_fallback"
8+
ModeBoth = "both"
9+
)
10+
11+
// NormalizeMode returns the effective fallback mode, defaulting to both.
12+
func NormalizeMode(mode string) string {
13+
if mode == "" {
14+
return ModeBoth
15+
}
16+
return mode
17+
}
18+
19+
// WriteEnabled reports whether fallback writes (write-through) are enabled.
20+
func WriteEnabled(mode string) bool {
21+
switch NormalizeMode(mode) {
22+
case ModeReadFallback:
23+
return false
24+
case ModeWriteThrough, ModeBoth:
25+
return true
26+
default:
27+
return true
28+
}
29+
}
30+
31+
// ValidateMode checks that mode is one of the supported values.
32+
func ValidateMode(mode string) error {
33+
if mode == "" {
34+
return nil
35+
}
36+
switch mode {
37+
case ModeWriteThrough, ModeReadFallback, ModeBoth:
38+
return nil
39+
default:
40+
return fmt.Errorf("fallback.mode must be %q, %q, or %q", ModeWriteThrough, ModeReadFallback, ModeBoth)
41+
}
42+
}

fallback/mode_test.go

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,25 @@
1+
package fallback
2+
3+
import (
4+
"testing"
5+
6+
"github.com/stretchr/testify/require"
7+
)
8+
9+
func TestNormalizeMode(t *testing.T) {
10+
require.Equal(t, ModeBoth, NormalizeMode(""))
11+
require.Equal(t, ModeReadFallback, NormalizeMode(ModeReadFallback))
12+
}
13+
14+
func TestWriteEnabled(t *testing.T) {
15+
require.True(t, WriteEnabled(ModeBoth))
16+
require.True(t, WriteEnabled(ModeWriteThrough))
17+
require.False(t, WriteEnabled(ModeReadFallback))
18+
require.True(t, WriteEnabled(""))
19+
}
20+
21+
func TestValidateMode(t *testing.T) {
22+
require.NoError(t, ValidateMode(""))
23+
require.NoError(t, ValidateMode(ModeBoth))
24+
require.Error(t, ValidateMode("invalid"))
25+
}

0 commit comments

Comments
 (0)