Skip to content

Commit b30d3b1

Browse files
Merge pull request #49 from chdb-io/fix/no-cancel-on-finished-stream
Do not cancel a stream the engine has already finished
2 parents 7188fcd + 2192fdb commit b30d3b1

1 file changed

Lines changed: 27 additions & 11 deletions

File tree

‎chdb-purego/streaming.go‎

Lines changed: 27 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,9 @@ type streamingResult struct {
66
curConn *chdb_connection
77
stream *chdb_result
88
curChunk ChdbResult
9+
// Set once the engine has signalled end of data, so Free() knows there is
10+
// nothing left to cancel. See the comment there.
11+
finished bool
912
}
1013

1114
func newStreamingResult(conn *chdb_connection, cRes *chdb_result) ChdbStreamResult {
@@ -37,7 +40,20 @@ func (c *streamingResult) Error() error {
3740
// Free implements ChdbStreamResult.
3841
func (c *streamingResult) Free() {
3942
if c.curConn != nil && c.stream != nil {
40-
chdbStreamCancelQuery(c.curConn, c.stream)
43+
// Cancel only while the engine still has a stream to stop. Once it has
44+
// signalled end of data it has retired the query's state, and cancelling a
45+
// retired stream walks into it: chdb_stream_cancel_query takes the handle as
46+
// an opaque pointer and casts it without checking, so there is nothing on
47+
// the engine side to turn that into an error instead of a fault. Against
48+
// chdb-core v26.7.0 it segfaults on linux/amd64, and survives on arm64 only
49+
// because the freed memory happens to still read back the way it did.
50+
//
51+
// Destroying is still correct and still required — the engine's own note on
52+
// cancel says the handle is released by chdb_destroy_query_result, not by
53+
// cancel — so only the cancel is conditional.
54+
if !c.finished {
55+
chdbStreamCancelQuery(c.curConn, c.stream)
56+
}
4157
chdbDestroyQueryResult(c.stream)
4258
}
4359

@@ -55,21 +71,21 @@ func (c *streamingResult) Cancel() {
5571

5672
// GetNext implements ChdbStreamResult.
5773
func (c *streamingResult) GetNext() ChdbResult {
58-
if c.curChunk == nil {
59-
nextChunk := chdbStreamFetchResult(c.curConn.internal_data, c.stream)
60-
if nextChunk == nil {
61-
return nil
62-
}
63-
c.curChunk = newChdbResult(nextChunk)
64-
return c.curChunk
74+
if c.curChunk != nil {
75+
// free the current chunk before getting the next one
76+
c.curChunk.Free()
77+
c.curChunk = nil
6578
}
66-
// free the current chunk before getting the next one
67-
c.curChunk.Free()
68-
c.curChunk = nil
6979
nextChunk := chdbStreamFetchResult(c.curConn.internal_data, c.stream)
7080
if nextChunk == nil {
81+
c.finished = true
7182
return nil
7283
}
7384
c.curChunk = newChdbResult(nextChunk)
85+
// A chunk with no rows is how the engine says the stream is done; callers
86+
// treat it as EOF and stop reading, so this is the last chunk there will be.
87+
if c.curChunk.RowsRead() == 0 {
88+
c.finished = true
89+
}
7490
return c.curChunk
7591
}

0 commit comments

Comments
 (0)