Skip to content

Commit cf7f877

Browse files
Async events (#190)
1 parent 74fb97e commit cf7f877

4 files changed

Lines changed: 71 additions & 66 deletions

File tree

‎sdk/adapters/supernodeservice/adapter.go‎

Lines changed: 29 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -446,6 +446,7 @@ func (a *cascadeAdapter) CascadeSupernodeDownload(
446446
bytesWritten int64
447447
chunkIndex int
448448
startedEmitted bool
449+
downloadStart time.Time
449450
)
450451

451452
// 3. Receive streamed responses
@@ -509,7 +510,11 @@ func (a *cascadeAdapter) CascadeSupernodeDownload(
509510
}
510511
}
511512
}
512-
in.EventLogger(ctx, toSdkEvent(x.Event.EventType), x.Event.Message, edata)
513+
// Avoid blocking Recv loop on event handling; dispatch asynchronously
514+
evtType := toSdkEvent(x.Event.EventType)
515+
go func(ed event.EventData, et event.EventType, msg string) {
516+
in.EventLogger(ctx, et, msg, ed)
517+
}(edata, evtType, x.Event.Message)
513518
}
514519

515520
// 3b. Actual data chunk
@@ -520,7 +525,10 @@ func (a *cascadeAdapter) CascadeSupernodeDownload(
520525
}
521526
if !startedEmitted {
522527
if in.EventLogger != nil {
523-
in.EventLogger(ctx, event.SDKDownloadStarted, "Download started", event.EventData{event.KeyActionID: in.ActionID})
528+
// mark start to compute throughput at completion
529+
downloadStart = time.Now()
530+
// Emit started asynchronously to avoid blocking
531+
go in.EventLogger(ctx, event.SDKDownloadStarted, "Download started", event.EventData{event.KeyActionID: in.ActionID})
524532
}
525533
startedEmitted = true
526534
}
@@ -538,7 +546,25 @@ func (a *cascadeAdapter) CascadeSupernodeDownload(
538546
a.logger.Info(ctx, "download complete", "bytes_written", bytesWritten, "path", in.OutputPath, "action_id", in.ActionID)
539547

540548
if in.EventLogger != nil {
541-
in.EventLogger(ctx, event.SDKDownloadCompleted, "Download completed", event.EventData{event.KeyActionID: in.ActionID, event.KeyOutputPath: in.OutputPath})
549+
// Compute metrics if we marked a start
550+
var elapsed float64
551+
var throughput float64
552+
if !downloadStart.IsZero() {
553+
elapsed = time.Since(downloadStart).Seconds()
554+
mb := float64(bytesWritten) / (1024.0 * 1024.0)
555+
if elapsed > 0 {
556+
throughput = mb / elapsed
557+
}
558+
}
559+
// Emit completion asynchronously with metrics
560+
go in.EventLogger(ctx, event.SDKDownloadCompleted, "Download completed", event.EventData{
561+
event.KeyActionID: in.ActionID,
562+
event.KeyOutputPath: in.OutputPath,
563+
event.KeyBytesTotal: bytesWritten,
564+
event.KeyChunks: chunkIndex,
565+
event.KeyElapsedSeconds: elapsed,
566+
event.KeyThroughputMBS: throughput,
567+
})
542568
}
543569
return &CascadeSupernodeDownloadResponse{
544570
Success: true,

‎sdk/net/factory.go‎

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -39,9 +39,10 @@ func NewClientFactory(ctx context.Context, logger log.Logger, keyring keyring.Ke
3939
// Tuned for 1GB max files with 4MB chunks
4040
// Reduce in-flight memory by aligning windows and msg sizes to chunk size.
4141
opts := client.DefaultClientOptions()
42-
opts.MaxRecvMsgSize = 8 * 1024 * 1024 // 8MB: supports 4MB chunks + overhead
43-
opts.MaxSendMsgSize = 8 * 1024 * 1024 // 8MB: supports 4MB chunks + overhead
44-
opts.InitialWindowSize = 4 * 1024 * 1024 // 4MB per-stream window ≈ chunk size
42+
opts.MaxRecvMsgSize = 12 * 1024 * 1024 // 8MB: supports 4MB chunks + overhead
43+
opts.MaxSendMsgSize = 12 * 1024 * 1024 // 8MB: supports 4MB chunks + overhead
44+
// Increase per-stream window to provide headroom for first data chunk + events
45+
opts.InitialWindowSize = 12 * 1024 * 1024 // 8MB per-stream window
4546
opts.InitialConnWindowSize = 64 * 1024 * 1024 // 64MB per-connection window
4647

4748
return &ClientFactory{

‎sdk/task/download.go‎

Lines changed: 10 additions & 56 deletions
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,6 @@ import (
44
"context"
55
"fmt"
66
"os"
7-
"sort"
87
"time"
98

109
"github.com/LumeraProtocol/supernode/v2/sdk/adapters/lumera"
@@ -77,51 +76,6 @@ func (t *CascadeDownloadTask) downloadFromSupernodes(ctx context.Context, supern
7776
}
7877
}
7978

80-
// Optionally rank supernodes by available memory to improve success for large files
81-
// We keep a short timeout per status fetch to avoid delaying downloads.
82-
type rankedSN struct {
83-
sn lumera.Supernode
84-
availGB float64
85-
hasStatus bool
86-
}
87-
ranked := make([]rankedSN, 0, len(supernodes))
88-
for _, sn := range supernodes {
89-
ranked = append(ranked, rankedSN{sn: sn})
90-
}
91-
92-
// Probe supernode status with short timeouts and close clients promptly
93-
for i := range ranked {
94-
sn := ranked[i].sn
95-
// 2s status timeout to keep this pass fast
96-
stx, cancel := context.WithTimeout(ctx, 2*time.Second)
97-
client, err := clientFactory.CreateClient(stx, sn)
98-
if err != nil {
99-
cancel()
100-
continue
101-
}
102-
status, err := client.GetSupernodeStatus(stx)
103-
_ = client.Close(stx)
104-
cancel()
105-
if err != nil {
106-
continue
107-
}
108-
ranked[i].hasStatus = true
109-
ranked[i].availGB = status.Resources.Memory.AvailableGB
110-
}
111-
112-
// Sort: nodes with status first, higher available memory first
113-
sort.Slice(ranked, func(i, j int) bool {
114-
if ranked[i].hasStatus != ranked[j].hasStatus {
115-
return ranked[i].hasStatus && !ranked[j].hasStatus
116-
}
117-
return ranked[i].availGB > ranked[j].availGB
118-
})
119-
120-
// Rebuild the supernodes list in the sorted order
121-
for i := range ranked {
122-
supernodes[i] = ranked[i].sn
123-
}
124-
12579
// Try supernodes sequentially, one by one (now sorted)
12680
var lastErr error
12781
for idx, sn := range supernodes {
@@ -146,8 +100,8 @@ func (t *CascadeDownloadTask) downloadFromSupernodes(ctx context.Context, supern
146100
continue
147101
}
148102

149-
// Success; return to caller
150-
return nil
103+
// Success; return to caller
104+
return nil
151105
}
152106

153107
if lastErr != nil {
@@ -176,15 +130,15 @@ func (t *CascadeDownloadTask) attemptDownload(
176130
t.LogEvent(ctx, evt, msg, data)
177131
}
178132

179-
resp, err := client.Download(ctx, req)
180-
if err != nil {
181-
return fmt.Errorf("download from %s: %w", sn.CosmosAddress, err)
182-
}
183-
if !resp.Success {
184-
return fmt.Errorf("download rejected by %s: %s", sn.CosmosAddress, resp.Message)
185-
}
133+
resp, err := client.Download(ctx, req)
134+
if err != nil {
135+
return fmt.Errorf("download from %s: %w", sn.CosmosAddress, err)
136+
}
137+
if !resp.Success {
138+
return fmt.Errorf("download rejected by %s: %s", sn.CosmosAddress, resp.Message)
139+
}
186140

187-
return nil
141+
return nil
188142
}
189143

190144
// downloadResult holds the result of a successful download attempt

‎supernode/node/action/server/cascade/cascade_action_server.go‎

Lines changed: 28 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -313,7 +313,14 @@ func (server *ActionServer) Download(req *pb.DownloadRequest, stream pb.CascadeS
313313
"chunk_size": chunkSize,
314314
})
315315

316-
// Announce: file is ready to be served to the client
316+
// Pre-read first chunk to avoid any delay between SERVE_READY and first data
317+
buf := make([]byte, chunkSize)
318+
n, readErr := f.Read(buf)
319+
if readErr != nil && readErr != io.EOF {
320+
return fmt.Errorf("chunked read failed: %w", readErr)
321+
}
322+
323+
// Announce: file is ready to be served to the client (right before first data)
317324
if err := stream.Send(&pb.DownloadResponse{
318325
ResponseType: &pb.DownloadResponse_Event{
319326
Event: &pb.DownloadEvent{
@@ -326,10 +333,27 @@ func (server *ActionServer) Download(req *pb.DownloadRequest, stream pb.CascadeS
326333
return err
327334
}
328335

329-
// Stream the file in fixed-size chunks
330-
buf := make([]byte, chunkSize)
336+
// Send pre-read first chunk if available
337+
if n > 0 {
338+
if err := stream.Send(&pb.DownloadResponse{
339+
ResponseType: &pb.DownloadResponse_Chunk{
340+
Chunk: &pb.DataChunk{Data: buf[:n]},
341+
},
342+
}); err != nil {
343+
logtrace.Error(ctx, "failed to stream first chunk", logtrace.Fields{logtrace.FieldError: err.Error()})
344+
return err
345+
}
346+
}
347+
348+
// If EOF after first read, we're done
349+
if readErr == io.EOF {
350+
logtrace.Info(ctx, "completed streaming all chunks", fields)
351+
return nil
352+
}
353+
354+
// Continue streaming remaining chunks
331355
for {
332-
n, readErr := f.Read(buf)
356+
n, readErr = f.Read(buf)
333357
if n > 0 {
334358
if err := stream.Send(&pb.DownloadResponse{
335359
ResponseType: &pb.DownloadResponse_Chunk{

0 commit comments

Comments
 (0)