-
Notifications
You must be signed in to change notification settings - Fork 2
Expand file tree
/
Copy pathformat.go
More file actions
112 lines (104 loc) · 3.29 KB
/
Copy pathformat.go
File metadata and controls
112 lines (104 loc) · 3.29 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
package main
import (
"context"
sppb "cloud.google.com/go/spanner/apiv1/spannerpb"
"github.com/apstndb/execspansql/jqresult"
"github.com/wader/gojq"
)
// writeResult chooses only an output strategy. Transaction selection, retries,
// and commit handling belong to executeQuery and runCLI.
func (c *preparedCommand) writeResult(ctx context.Context, result *queryResult, sinks *outputSinks) error {
// Discarded results need metadata/stats only. DML is already materialized;
// read-only CSV and lazy jq retain their streaming behavior.
if c.DiscardResults || result.resultSet != nil || (c.Format != "experimental_csv" && c.jqMode == jqresult.InputEager) {
rs, err := result.materialize(materializeWithoutRows(c.opts))
if err != nil {
return err
}
return c.writeResultSet(ctx, rs, sinks)
}
if c.Format == "experimental_csv" {
result, err := writeCsvFromRowIter(sinks.primary, result.rowIter, c.RedactRows)
if err != nil {
return err
}
sinks.MarkPrimaryComplete()
if !sinks.hasPlan {
return nil
}
stats, err := statsFromWriterResult(result)
if err != nil {
return err
}
return c.writePlan(ctx, sinks, result.Metadata, stats)
}
enc, err := newEncoder(sinks.primary, c.Format, c.CompactOutput, c.JqRawOutput)
if err != nil {
return err
}
var lazyOpts []jqresult.LazyOption
if sinks.hasPlan {
lazyOpts = append(lazyOpts, jqresult.WithOmitQueryPlan())
}
lazy := jqresult.NewLazy(result.rowIter, c.RedactRows, lazyOpts...)
defer lazy.Stop()
if err := printJQ(ctx, c.jqCode, lazy, enc); err != nil {
return err
}
sinks.MarkPrimaryComplete()
if !sinks.hasPlan {
return nil
}
if err := lazy.Drain(); err != nil {
return err
}
drained := lazy.Result()
stats, err := drained.StatsProto()
if err != nil {
return err
}
return c.writePlan(ctx, sinks, drained.Metadata, stats)
}
func (c *preparedCommand) writeResultSet(ctx context.Context, rs *sppb.ResultSet, sinks *outputSinks) error {
stats, metadata := rs.Stats, rs.Metadata
if sinks.hasPlan {
stats, metadata = stripQueryPlanForPrimary(rs)
}
if sinks.primary != nil {
if c.Format == "experimental_csv" {
if err := writeCsvFromResultSet(sinks.primary, rs); err != nil {
return err
}
} else {
input, err := jqresult.ResultSetMap(rs)
if err != nil {
return err
}
enc, err := newEncoder(sinks.primary, c.Format, c.CompactOutput, c.JqRawOutput)
if err != nil {
return err
}
if err := printJQ(ctx, c.jqCode, input, enc); err != nil {
return err
}
}
}
sinks.MarkPrimaryComplete()
return c.writePlan(ctx, sinks, metadata, stats)
}
// printJQ owns encoder completion on both success and failure. Use the caller's
// cancellation context even after SQL has completed and no RPC remains active.
func printJQ(ctx context.Context, code *gojq.Code, input any, enc encoder) (err error) {
defer func() {
if closeErr := closeEncoder(enc); err == nil {
err = closeErr
}
}()
return jqresult.Print(enc, code.RunWithContext(ctx, input))
}
// Use the prepared statement for query text so rendering never reloads SQL files.
func (c *preparedCommand) writePlan(ctx context.Context, sinks *outputSinks, metadata *sppb.ResultSetMetadata, stats *sppb.ResultSetStats) error {
o := c.opts
o.Sql = c.statement.SQL
return writePlan(ctx, sinks.plan, effectivePlanFormat(o), metadata, stats, o)
}