@@ -37,6 +37,8 @@ import (
3737 "github.com/transparency-dev/tessera/storage/posix"
3838 "golang.org/x/net/http2"
3939
40+ sdbNote "golang.org/x/mod/sumdb/note"
41+
4042 "log/slog"
4143)
4244
@@ -47,6 +49,9 @@ func init() {
4749var (
4850 mirrorURL multiStringFlag
4951
52+ sourceLogURL = flag .String ("source_log_url" , "" , "URL of the log to mirror, or empty if using a locally created log." )
53+ sourceLogPubKeyPath = flag .String ("source_log_pubkey_path" , "" , "Path to the public key of the log to mirror. Required if --source_log_url is set." )
54+
5055 storageDir = flag .String ("storage_dir" , "" , "Root directory to store log data" )
5156 logPrivKey = flag .String ("log_private_key" , "" , "Location of private key file" )
5257
7176 hc * http.Client
7277)
7378
79+ type logReader interface {
80+ ReadCheckpoint (ctx context.Context ) ([]byte , error )
81+ ReadEntryBundle (ctx context.Context , index uint64 , p uint8 ) ([]byte , error )
82+ ReadTile (ctx context.Context , level , index uint64 , p uint8 ) ([]byte , error )
83+ }
84+
7485func main () {
7586 flag .Parse ()
7687 slog .SetDefault (slog .New (slog .NewTextHandler (os .Stderr , & slog.HandlerOptions {Level : slog .Level (* slogLevel )})))
@@ -102,42 +113,15 @@ func main() {
102113
103114 ctx , cancel := context .WithCancel (context .Background ())
104115
105- // Gather the info needed for reading/writing checkpoints.
106- s := getSignerOrDie (ctx )
107-
108- // Create the Tessera POSIX storage, using the directory from the --storage_dir flag.
109- driver , err := posix .New (ctx , posix.Config {Path : * storageDir })
110- if err != nil {
111- slog .ErrorContext (ctx , "Failed to construct storage" , slog .Any ("error" , err ))
112- os .Exit (1 )
113- }
114-
115- appender , shutdown , logReader , err := tessera .NewAppender (ctx , driver , tessera .NewAppendOptions ().
116- WithCheckpointSigner (s ).
117- WithCheckpointInterval (time .Second ).
118- WithCheckpointRepublishInterval (time .Minute ).
119- WithBatching (tessera .DefaultBatchMaxSize , time .Second ))
120- if err != nil {
121- slog .ErrorContext (ctx , "Failed to create new appender" , slog .Any ("error" , err ))
122- os .Exit (1 )
123- }
116+ appender , logReader , shutdown , verifier := logFromFlags (ctx )
124117 defer func () {
125118 if err := shutdown (ctx ); err != nil {
126119 slog .ErrorContext (ctx , "Failed to shutdown appender" , slog .Any ("error" , err ))
127120 }
128121 }()
129122
130- leafWriter := func (ctx context.Context , newLeaf []byte ) (uint64 , error ) {
131- idx , err := appender .Add (ctx , tessera .NewEntry (newLeaf ))()
132- if err != nil {
133- return 0 , fmt .Errorf ("failed to add leaf: %v" , err )
134- }
135- return idx .Index , nil
136- }
137-
138- var cpRaw []byte
139123 cons := client .UnilateralConsensus (logReader .ReadCheckpoint )
140- tracker , err := client .NewLogStateTracker (ctx , logReader .ReadTile , cpRaw , s . Verifier (), s . Verifier () .Name (), cons )
124+ tracker , err := client .NewLogStateTracker (ctx , logReader .ReadTile , nil , verifier , verifier .Name (), cons )
141125 if err != nil {
142126 slog .ErrorContext (ctx , "Failed to create LogStateTracker" , slog .Any ("error" , err ))
143127 os .Exit (1 )
@@ -160,7 +144,7 @@ func main() {
160144 mOpts := mirror .NewOptions ().
161145 WithMirrorURL (mURL ).
162146 WithHTTPClient (hc ).
163- WithLogOrigin (s . Verifier () .Name ()).
147+ WithLogOrigin (verifier .Name ()).
164148 WithTileFetcher (logReader .ReadTile ).
165149 WithBundleFetcher (logReader .ReadEntryBundle ).
166150 WithMirrorCheckpointFetcher (func (ctx context.Context ) ([]byte , error ) {
@@ -175,18 +159,13 @@ func main() {
175159 go runMirrorSync (ctx , tracker , mc , mURL )
176160 }
177161
178- ha := loadtest .NewHammerAnalyser (func () uint64 { return tracker .Latest ().Size })
179- ha .Run (ctx )
180-
181- gen := newLeafGenerator (tracker .Latest ().Size , * leafMinSize , * dupChance )
182- opts := loadtest.HammerOpts {
183- MaxReadOpsPerSecond : 0 ,
184- MaxWriteOpsPerSecond : * maxWriteOpsPerSecond ,
185- NumReadersRandom : 0 ,
186- NumReadersFull : 0 ,
187- NumWriters : * numWriters ,
162+ var hammer * loadtest.Hammer
163+ var ha * loadtest.HammerAnalyser
164+ if appender != nil {
165+ hammer , ha = newHammer (ctx , appender , logReader , tracker )
166+ ha .Run (ctx )
167+ hammer .Run (ctx )
188168 }
189- hammer := loadtest .NewHammer (tracker , logReader .ReadEntryBundle , leafWriter , gen , ha .SeqLeafChan , ha .ErrChan , opts )
190169
191170 exitCode := 0
192171 if * leafWriteGoal > 0 {
@@ -226,10 +205,9 @@ func main() {
226205 }
227206 }()
228207 }
229- hammer .Run (ctx )
230208
231209 if * showUI {
232- c := loadtest .NewController (hammer , ha )
210+ c := loadtest .NewController (hammer , ha , tracker )
233211 c .Run (ctx )
234212 } else {
235213 <- ctx .Done ()
@@ -350,3 +328,100 @@ func runMirrorSync(ctx context.Context, tracker *client.LogStateTracker, mClient
350328 }
351329 }
352330}
331+
332+ func logFromFlags (ctx context.Context ) (* tessera.Appender , logReader , func (context.Context ) error , sdbNote.Verifier ) {
333+ if * sourceLogURL != "" {
334+ if * sourceLogPubKeyPath == "" {
335+ slog .ErrorContext (ctx , "Source log URL provided but no public key path." )
336+ os .Exit (1 )
337+ }
338+ lr , shutdown , v , err := newRemoteLog (ctx , * sourceLogURL , * sourceLogPubKeyPath )
339+ if err != nil {
340+ slog .ErrorContext (ctx , "Failed to create remote log" , slog .Any ("error" , err ))
341+ os .Exit (1 )
342+ }
343+ return nil , lr , shutdown , v
344+ }
345+
346+ if * storageDir == "" {
347+ slog .ErrorContext (ctx , "No storage directory provided." )
348+ os .Exit (1 )
349+ }
350+ app , lr , shutdown , v , err := newLocalLog (ctx , * storageDir , getSignerOrDie (ctx ))
351+ if err != nil {
352+ slog .ErrorContext (ctx , "Failed to create local log" , slog .Any ("error" , err ))
353+ os .Exit (1 )
354+ }
355+ return app , lr , shutdown , v
356+ }
357+
358+ // newLocalLog creates a new local POSIX log which can be used by a hammer to hammer the mirror.
359+ //
360+ // Returns an Appender for adding leaves to the log, a LogReader that reads from the local log,
361+ func newLocalLog (ctx context.Context , storageDir string , s note.Signer ) (* tessera.Appender , logReader , func (context.Context ) error , sdbNote.Verifier , error ) {
362+ // Create the Tessera POSIX storage, using the provided directory.
363+ driver , err := posix .New (ctx , posix.Config {Path : storageDir })
364+ if err != nil {
365+ return nil , nil , nil , nil , fmt .Errorf ("failed to construct storage: %w" , err )
366+ }
367+
368+ appender , shutdown , logReader , err := tessera .NewAppender (ctx , driver , tessera .NewAppendOptions ().
369+ WithCheckpointSigner (s ).
370+ WithCheckpointInterval (time .Second ).
371+ WithCheckpointRepublishInterval (time .Minute ).
372+ WithBatching (tessera .DefaultBatchMaxSize , time .Second ))
373+ if err != nil {
374+ return nil , nil , nil , nil , err
375+ }
376+ return appender , logReader , shutdown , s .Verifier (), nil
377+ }
378+
379+ // newRemoteLog creates a logReader that reads from a remote log, and a shutdown func.
380+ func newRemoteLog (_ context.Context , logURL string , pubKeyPath string ) (logReader , func (context.Context ) error , sdbNote.Verifier , error ) {
381+ pubKRaw , err := getKeyFile (pubKeyPath )
382+ if err != nil {
383+ return nil , nil , nil , fmt .Errorf ("failed to read public key %q: %v" , pubKeyPath , err )
384+ }
385+
386+ pubKey , err := note .NewVerifier (pubKRaw )
387+ if err != nil {
388+ return nil , nil , nil , fmt .Errorf ("failed to parse public key %q: %v" , pubKeyPath , err )
389+ }
390+
391+ var lr logReader
392+ switch logURL := strings .ToLower (logURL ); {
393+ case strings .HasPrefix (logURL , "http://" ), strings .HasPrefix (logURL , "https://" ):
394+ u , err := url .Parse (logURL )
395+ if err != nil {
396+ return nil , nil , nil , fmt .Errorf ("failed to parse log URL %q: %v" , logURL , err )
397+ }
398+ lr , err = client .NewHTTPFetcher (u , hc )
399+ default :
400+ lr = client.FileFetcher {Root : logURL }
401+ }
402+
403+ return lr , func (context.Context ) error { return nil }, pubKey , nil
404+ }
405+
406+ func newHammer (_ context.Context , appender * tessera.Appender , logReader logReader , tracker * client.LogStateTracker ) (* loadtest.Hammer , * loadtest.HammerAnalyser ) {
407+ leafWriter := func (ctx context.Context , newLeaf []byte ) (uint64 , error ) {
408+ idx , err := appender .Add (ctx , tessera .NewEntry (newLeaf ))()
409+ if err != nil {
410+ return 0 , fmt .Errorf ("failed to add leaf: %v" , err )
411+ }
412+ return idx .Index , nil
413+ }
414+ ha := loadtest .NewHammerAnalyser (func () uint64 { return tracker .Latest ().Size })
415+
416+ gen := newLeafGenerator (tracker .Latest ().Size , * leafMinSize , * dupChance )
417+ opts := loadtest.HammerOpts {
418+ MaxReadOpsPerSecond : 0 ,
419+ MaxWriteOpsPerSecond : * maxWriteOpsPerSecond ,
420+ NumReadersRandom : 0 ,
421+ NumReadersFull : 0 ,
422+ NumWriters : * numWriters ,
423+ }
424+ hammer := loadtest .NewHammer (tracker , logReader .ReadEntryBundle , leafWriter , gen , ha .SeqLeafChan , ha .ErrChan , opts )
425+
426+ return hammer , ha
427+ }
0 commit comments