| 1 | // SPDX-License-Identifier: GPL-3.0-or-later |
| 2 | |
| 3 | package functions |
| 4 | |
| 5 | import ( |
| 6 | "bytes" |
| 7 | "encoding/csv" |
| 8 | "errors" |
| 9 | "fmt" |
| 10 | "strconv" |
| 11 | "strings" |
| 12 | "time" |
| 13 | ) |
| 14 | |
| 15 | const ( |
| 16 | lineFunction = "FUNCTION" |
| 17 | lineFunctionPayload = "FUNCTION_PAYLOAD" |
| 18 | lineFunctionPayloadEnd = "FUNCTION_PAYLOAD_END" |
| 19 | lineFunctionCancel = "FUNCTION_CANCEL" |
| 20 | lineFunctionProgress = "FUNCTION_PROGRESS" |
| 21 | lineQuit = "QUIT" |
| 22 | ) |
| 23 | |
| 24 | type Function struct { |
| 25 | key string |
| 26 | UID string |
| 27 | Timeout time.Duration |
| 28 | Name string |
| 29 | Args []string |
| 30 | Payload []byte |
| 31 | Permissions string |
| 32 | Source string |
| 33 | ContentType string |
| 34 | } |
| 35 | |
| 36 | func (f *Function) String() string { |
| 37 | return fmt.Sprintf("key: '%s', uid: '%s', timeout: '%s', function: '%s', args: '%v', permissions: '%s', source: '%s', contentType: '%s', payload: '%s'", |
| 38 | f.key, f.UID, f.Timeout, f.Name, f.Args, f.Permissions, f.Source, f.ContentType, string(f.Payload)) |
| 39 | } |
| 40 | |
| 41 | func newInputParser() *inputParser { |
| 42 | return &inputParser{} |
| 43 | } |
| 44 | |
| 45 | type inputParser struct { |
| 46 | currentFn *Function |
| 47 | readingPayload bool |
| 48 | payloadBuf bytes.Buffer |
| 49 | } |
| 50 | |
| 51 | func (p *inputParser) parse(line string) (*Function, error) { |
| 52 | event, err := p.parseEvent(line) |
| 53 | if err != nil { |
| 54 | return nil, err |
| 55 | } |
| 56 | if event.kind == inputEventCall { |
| 57 | return event.fn, nil |
| 58 | } |
| 59 | return nil, nil |
| 60 | } |
| 61 | |
| 62 | type inputEventKind uint8 |
| 63 | |
| 64 | const ( |
| 65 | inputEventNone inputEventKind = iota |
| 66 | inputEventCall |
| 67 | inputEventCancel |
| 68 | inputEventProgress |
| 69 | inputEventQuit |
| 70 | ) |
| 71 | |
| 72 | type inputEvent struct { |
| 73 | kind inputEventKind |
| 74 | fn *Function |
| 75 | uid string |
| 76 | preAdmission bool |
| 77 | } |
| 78 | |
| 79 | func (p *inputParser) parseEvent(line string) (inputEvent, error) { |
| 80 | if line = strings.TrimSpace(line); line == "" { |
| 81 | return inputEvent{}, nil |
| 82 | } |
| 83 | |
| 84 | if p.readingPayload { |
| 85 | return p.handlePayloadLine(line) |
| 86 | } |
| 87 | |
| 88 | switch { |
| 89 | case line == lineQuit: |
| 90 | return inputEvent{kind: inputEventQuit}, nil |
| 91 | case hasLinePrefix(line, lineFunctionCancel): |
| 92 | return parseCancelEvent(line) |
| 93 | case hasLinePrefix(line, lineFunctionProgress): |
| 94 | return parseProgressEvent(line), nil |
| 95 | case strings.HasPrefix(line, lineFunction+" "): |
| 96 | fn, err := p.parseFunction(line) |
| 97 | if err != nil { |
| 98 | return inputEvent{}, err |
| 99 | } |
| 100 | return inputEvent{kind: inputEventCall, fn: fn}, nil |
| 101 | case strings.HasPrefix(line, lineFunctionPayload+" "): |
| 102 | fn, err := p.parseFunction(line) |
| 103 | if err != nil { |
| 104 | return inputEvent{}, err |
| 105 | } |
| 106 | p.readingPayload = true |
| 107 | p.currentFn = fn |
| 108 | p.payloadBuf.Reset() |
| 109 | return inputEvent{}, nil |
| 110 | default: |
| 111 | return inputEvent{}, errors.New("unexpected line format") |
| 112 | } |
| 113 | } |
| 114 | |
| 115 | func (p *inputParser) handlePayloadLine(line string) (inputEvent, error) { |
| 116 | if line == lineFunctionPayloadEnd { |
| 117 | p.readingPayload = false |
| 118 | p.currentFn.Payload = []byte(p.payloadBuf.String()) |
| 119 | fn := p.currentFn |
| 120 | p.currentFn = nil |
| 121 | p.payloadBuf.Reset() |
| 122 | return inputEvent{kind: inputEventCall, fn: fn}, nil |
| 123 | } |
| 124 | |
| 125 | if hasLinePrefix(line, lineFunctionCancel) { |
| 126 | event, err := parseCancelEvent(line) |
| 127 | if err != nil { |
| 128 | // Malformed cancel must not affect payload parser state. |
| 129 | return inputEvent{}, err |
| 130 | } |
| 131 | |
| 132 | if p.currentFn != nil && event.uid == p.currentFn.UID { |
| 133 | p.resetPayloadState() |
| 134 | event.preAdmission = true |
| 135 | } |
| 136 | return event, nil |
| 137 | } |
| 138 | |
| 139 | if hasLinePrefix(line, lineFunctionProgress) { |
| 140 | return parseProgressEvent(line), nil |
| 141 | } |
| 142 | |
| 143 | if line == lineQuit { |
| 144 | p.resetPayloadState() |
| 145 | return inputEvent{kind: inputEventQuit}, nil |
| 146 | } |
| 147 | |
| 148 | if hasLinePrefix(line, lineFunction) || strings.HasPrefix(line, lineFunction+"_") { |
| 149 | p.resetPayloadState() |
| 150 | return p.parseEvent(line) |
| 151 | } |
| 152 | |
| 153 | if p.payloadBuf.Len() > 0 { |
| 154 | p.payloadBuf.WriteByte('\n') |
| 155 | } |
| 156 | p.payloadBuf.WriteString(line) |
| 157 | |
| 158 | return inputEvent{}, nil |
| 159 | } |
| 160 | |
| 161 | func (p *inputParser) resetPayloadState() { |
| 162 | p.readingPayload = false |
| 163 | p.currentFn = nil |
| 164 | p.payloadBuf.Reset() |
| 165 | } |
| 166 | |
| 167 | func hasLinePrefix(line, keyword string) bool { |
| 168 | return line == keyword || strings.HasPrefix(line, keyword+" ") |
| 169 | } |
| 170 | |
| 171 | func parseCancelEvent(line string) (inputEvent, error) { |
| 172 | parts := strings.Fields(line) |
| 173 | if len(parts) != 2 || parts[0] != lineFunctionCancel || parts[1] == "" { |
| 174 | return inputEvent{}, errors.New("unexpected FUNCTION_CANCEL format") |
| 175 | } |
| 176 | return inputEvent{kind: inputEventCancel, uid: parts[1]}, nil |
| 177 | } |
| 178 | |
| 179 | func parseProgressEvent(line string) inputEvent { |
| 180 | parts := strings.Fields(line) |
| 181 | event := inputEvent{kind: inputEventProgress} |
| 182 | if len(parts) >= 2 { |
| 183 | event.uid = parts[1] |
| 184 | } |
| 185 | return event |
| 186 | } |
| 187 | |
| 188 | func (p *inputParser) parseFunction(line string) (*Function, error) { |
| 189 | r := csv.NewReader(strings.NewReader(line)) |
| 190 | r.Comma = ' ' |
| 191 | |
| 192 | parts, err := r.Read() |
| 193 | if err != nil { |
| 194 | return nil, fmt.Errorf("failed to parse CSV: %w", err) |
| 195 | } |
| 196 | |
| 197 | if n := len(parts); n != 6 && n != 7 { |
| 198 | return nil, fmt.Errorf("unexpected number of parts: want 6 or 7, got %d", n) |
| 199 | } |
| 200 | |
| 201 | timeout, err := strconv.ParseInt(parts[2], 10, 64) |
| 202 | if err != nil { |
| 203 | return nil, fmt.Errorf("invalid timeout value: %w", err) |
| 204 | } |
| 205 | |
| 206 | nameAndArgs := strings.Split(parts[3], " ") |
| 207 | if len(nameAndArgs) == 0 { |
| 208 | return nil, fmt.Errorf("empty function name and arguments") |
| 209 | } |
| 210 | |
| 211 | fn := &Function{ |
| 212 | key: parts[0], |
| 213 | UID: parts[1], |
| 214 | Timeout: time.Duration(timeout) * time.Second, |
| 215 | Name: nameAndArgs[0], |
| 216 | Args: nameAndArgs[1:], |
| 217 | Permissions: parts[4], |
| 218 | Source: parts[5], |
| 219 | } |
| 220 | |
| 221 | if len(parts) == 7 { |
| 222 | fn.ContentType = parts[6] |
| 223 | } |
| 224 | |
| 225 | return fn, nil |
| 226 | } |