A11 (C++ runtime)
Native C++ implementation of the A11 streaming action runtime
Loading...
Searching...
No Matches
redis_chunk_store_script.h
Go to the documentation of this file.
1/*
2 * Copyright 2026 The A11 Authors
3 *
4 * Licensed under the Apache License, Version 2.0 (the "License");
5 * you may not use this file except in compliance with the License.
6 * You may obtain a copy of the License at
7 *
8 * http://www.apache.org/licenses/LICENSE-2.0
9 *
10 * Unless required by applicable law or agreed to in writing, software
11 * distributed under the License is distributed on an "AS IS" BASIS,
12 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13 * See the License for the specific language governing permissions and
14 * limitations under the License.
15 */
16
17#ifndef A11_STORES_REDIS_CHUNK_STORE_SCRIPT_H_
18#define A11_STORES_REDIS_CHUNK_STORE_SCRIPT_H_
19
20#include <string_view>
21
23
24// One dispatcher keeps every operation on the same versioned schema and the
25// same explicitly declared key set. Redis Cluster can therefore verify that
26// all touched keys share the caller-selected node slot.
27inline constexpr std::string_view kRedisChunkStoreScript = R"lua(
28local meta_key = KEYS[1]
29local stream_key = KEYS[2]
30local seq_key = KEYS[3]
31local arrival_key = KEYS[4]
32local blobs_key = KEYS[5]
33local events_channel = KEYS[6]
34
35local operation = ARGV[1]
36local node_id = ARGV[2]
37local max_seq_value = 4294967295
38local max_count_value = 4294967296
39-- At most 2^32 successful puts, 2^32 first-time clears, and one close can
40-- publish a revision. Keeping every Lua counter below 2^53 also avoids the
41-- lossy integer conversions of Redis's Lua 5.1 number type.
42local max_revision_value = 8589934593
43
44if type(operation) ~= 'string' or type(node_id) ~= 'string' or node_id == '' then
45 return {'error', 'INVALID_ARGUMENT', 'Invalid chunk-store script envelope'}
46end
47
48local function failure(code, message)
49 return {'error', code, message}
50end
51
52local function key_type(key)
53 local reply = redis.call('TYPE', key)
54 if type(reply) == 'table' then
55 return reply['ok']
56 end
57 return reply
58end
59
60local function check_type(key, expected)
61 local actual = key_type(key)
62 return actual == 'none' or actual == expected
63end
64
65local function is_canonical_decimal(value, maximum)
66 if type(value) ~= 'string' or value == '' or
67 string.find(value, '[^0-9]') then
68 return false
69 end
70 if #value > 1 and string.sub(value, 1, 1) == '0' then
71 return false
72 end
73 if #value ~= #maximum then
74 return #value < #maximum
75 end
76 return value <= maximum
77end
78
79local function canonical_uint(value, maximum, maximum_text)
80 if not is_canonical_decimal(value, maximum_text) then
81 return nil
82 end
83 local parsed = tonumber(value)
84 if not parsed or parsed < 0 or parsed > maximum or
85 math.floor(parsed) ~= parsed then
86 return nil
87 end
88 return parsed
89end
90
91if not check_type(meta_key, 'hash') then
92 return failure('DATA_LOSS', 'Chunk store metadata key is not a hash')
93end
94if not check_type(stream_key, 'stream') then
95 return failure('DATA_LOSS', 'Chunk store data key is not a stream')
96end
97if not check_type(seq_key, 'hash') then
98 return failure('DATA_LOSS', 'Chunk store sequence index is not a hash')
99end
100if not check_type(arrival_key, 'hash') then
101 return failure('DATA_LOSS', 'Chunk store arrival index is not a hash')
102end
103if not check_type(blobs_key, 'hash') then
104 return failure('DATA_LOSS', 'Chunk store blob key is not a hash')
105end
106
107local meta_exists = redis.call('EXISTS', meta_key) == 1
108local metadata_closed = false
109local metadata_size = 0
110local metadata_put_count = 0
111local metadata_next_cursor = 0
112local metadata_revision = 0
113local metadata_final_seq = nil
114local metadata_max_seq = nil
115if not meta_exists then
116 if redis.call('EXISTS', stream_key, seq_key, arrival_key, blobs_key) ~= 0 then
117 return failure('DATA_LOSS', 'Chunk data exists without store metadata')
118 end
119else
120 local stored_id = redis.call('HGET', meta_key, 'id')
121 local schema = redis.call('HGET', meta_key, 'schema')
122 local closed = redis.call('HGET', meta_key, 'closed')
123 local size = redis.call('HGET', meta_key, 'size')
124 local put_count = redis.call('HGET', meta_key, 'put_count')
125 local next_cursor = redis.call('HGET', meta_key, 'next_cursor')
126 local revision = redis.call('HGET', meta_key, 'revision')
127 local final_seq = redis.call('HGET', meta_key, 'final_seq')
128 local max_seq = redis.call('HGET', meta_key, 'max_seq')
129 metadata_size = canonical_uint(size, max_count_value, '4294967296')
130 metadata_put_count = canonical_uint(
131 put_count, max_count_value, '4294967296')
132 metadata_next_cursor = canonical_uint(
133 next_cursor, max_count_value, '4294967296')
134 metadata_revision = canonical_uint(
135 revision, max_revision_value, '8589934593')
136 if final_seq then
137 metadata_final_seq = canonical_uint(
138 final_seq, max_seq_value, '4294967295')
139 end
140 if max_seq then
141 metadata_max_seq = canonical_uint(max_seq, max_seq_value, '4294967295')
142 end
143 if not stored_id or schema ~= '1' or
144 (closed ~= '0' and closed ~= '1') or metadata_size == nil or
145 metadata_put_count == nil or metadata_next_cursor == nil or
146 metadata_revision == nil or
147 (final_seq and metadata_final_seq == nil) or
148 (max_seq and metadata_max_seq == nil) then
149 return failure('DATA_LOSS', 'Chunk store metadata is incomplete or corrupt')
150 end
151 if stored_id ~= node_id then
152 return failure('DATA_LOSS', 'Chunk store key belongs to a different node')
153 end
154 metadata_closed = closed == '1'
155 local stored_status = redis.call('HGET', meta_key, 'status')
156 if metadata_closed and not stored_status then
157 return failure('DATA_LOSS', 'Closed chunk store has no terminal status')
158 end
159 if not metadata_closed and stored_status then
160 return failure('DATA_LOSS', 'Open chunk store has a terminal status')
161 end
162 if metadata_size ~= metadata_put_count or
163 metadata_next_cursor > metadata_size then
164 return failure('DATA_LOSS', 'Chunk store counters are inconsistent')
165 end
166 if (metadata_size == 0 and metadata_max_seq ~= nil) or
167 (metadata_size > 0 and metadata_max_seq == nil) or
168 (metadata_final_seq ~= nil and
169 metadata_final_seq ~= metadata_max_seq) then
170 return failure('DATA_LOSS', 'Chunk store sequence metadata is inconsistent')
171 end
172 if redis.call('HLEN', seq_key) ~= metadata_size or
173 redis.call('HLEN', arrival_key) ~= metadata_put_count then
174 return failure('DATA_LOSS', 'Chunk store index cardinality is inconsistent')
175 end
176 if redis.call('HLEN', blobs_key) > metadata_size then
177 return failure('DATA_LOSS', 'Chunk store contains unindexed blobs')
178 end
179 local expected_stream_size = metadata_size + (metadata_closed and 1 or 0)
180 if redis.call('XLEN', stream_key) ~= expected_stream_size then
181 return failure('DATA_LOSS', 'Chunk store stream cardinality is inconsistent')
182 end
183end
184
185local function ensure_meta()
186 if not meta_exists then
187 redis.call('HSET', meta_key,
188 'schema', '1',
189 'id', node_id,
190 'closed', '0',
191 'size', '0',
192 'put_count', '0',
193 'next_cursor', '0',
194 'revision', '0')
195 meta_exists = true
196 end
197end
198
199local function revision()
200 if not meta_exists then
201 return '0'
202 end
203 return redis.call('HGET', meta_key, 'revision') or '0'
204end
205
206local function publish_change()
207 -- Every caller checks this before its first mutation. HINCRBY therefore
208 -- cannot overflow after a batch has partially committed.
209 local changed = redis.call('HINCRBY', meta_key, 'revision', 1)
210 redis.call('PUBLISH', events_channel, tostring(changed))
211 return changed
212end
213
214local function can_publish_change()
215 if metadata_revision >= max_revision_value then
216 return failure('RESOURCE_EXHAUSTED',
217 'Chunk store mutation revision exhausted')
218 end
219 return nil
220end
221
222local function field_map(entry)
223 if #entry == 0 then
224 return nil
225 end
226 local fields = entry[1][2]
227 local result = {}
228 for index = 1, #fields, 2 do
229 result[fields[index]] = fields[index + 1]
230 end
231 return result
232end
233
234local function can_append_stream_entry()
235 if redis.call('EXISTS', stream_key) == 0 then
236 return true
237 end
238 local info = redis.call('XINFO', 'STREAM', stream_key)
239 local stream_id = nil
240 for index = 1, #info, 2 do
241 if info[index] == 'last-generated-id' then
242 stream_id = info[index + 1]
243 break
244 end
245 end
246 if not stream_id then
247 return false
248 end
249 local separator = string.find(stream_id, '-', 1, true)
250 if not separator then
251 return false
252 end
253 -- XADD * may increment the sequence component, but no automatic ID can be
254 -- greater once the millisecond component itself has reached uint64_t max.
255 return string.sub(stream_id, 1, separator - 1) ~= '18446744073709551615'
256end
257
258-- Returns state, payload-or-ref, stream id, arrival, and storage. Every
259-- payload is an encoded Chunk. Tombstones retain a separately prepared,
260-- data-free encoding so ClearData can preserve metadata atomically without
261-- teaching Lua A11's MessagePack schema.
262local function load_chunk(seq)
263 local stream_id = redis.call('HGET', seq_key, tostring(seq))
264 if not stream_id then
265 return 'missing', '', '', '', ''
266 end
267 local entry = redis.call('XRANGE', stream_key, stream_id, stream_id, 'COUNT', 1)
268 local fields = field_map(entry)
269 if not fields then
270 return 'data_loss', 'Sequence index references a missing stream entry',
271 stream_id, '', ''
272 end
273 if fields['v'] ~= '1' or fields['kind'] ~= 'chunk' or
274 fields['seq'] ~= tostring(seq) or not fields['arrival'] then
275 return 'data_loss', 'Stream entry metadata is corrupt', stream_id, '', ''
276 end
277 local arrival = canonical_uint(
278 fields['arrival'], max_seq_value, '4294967295')
279 if arrival == nil or
280 redis.call('HGET', arrival_key, fields['arrival']) ~= tostring(seq) then
281 return 'data_loss', 'Stream entry arrival index is corrupt', stream_id,
282 fields['arrival'], ''
283 end
284 local storage = fields['storage']
285 if storage == 'inline' then
286 if fields['payload'] == nil then
287 return 'data_loss', 'Inline stream entry has no payload', stream_id,
288 fields['arrival'], storage
289 end
290 return 'item', fields['payload'], stream_id, fields['arrival'], storage
291 end
292 if storage == 'redis' then
293 if fields['ref'] ~= tostring(seq) then
294 return 'data_loss', 'Redis stream entry has no blob reference', stream_id,
295 fields['arrival'], storage
296 end
297 local payload = redis.call('HGET', blobs_key, fields['ref'])
298 if payload == false then
299 return 'data_loss', 'Redis stream entry references a missing blob',
300 stream_id, fields['arrival'], storage
301 end
302 return 'item', payload, stream_id, fields['arrival'], storage
303 end
304 if storage == 'tombstone' then
305 if fields['payload'] == nil then
306 return 'data_loss', 'Tombstone stream entry has no payload', stream_id,
307 fields['arrival'], storage
308 end
309 return 'item', fields['payload'], stream_id, fields['arrival'], storage
310 end
311 if storage == 's3' then
312 return 's3', fields['ref'] or '', stream_id, fields['arrival'], storage
313 end
314 return 'data_loss', 'Stream entry has an unknown storage kind', stream_id,
315 fields['arrival'], storage or ''
316end
317
318local function final_seq()
319 if not meta_exists then
320 return ''
321 end
322 return redis.call('HGET', meta_key, 'final_seq') or ''
323end
324
325local function missing_result(kind, value)
326 if meta_exists and redis.call('HGET', meta_key, 'closed') == '1' then
327 return {'closed', redis.call('HGET', meta_key, 'status')}
328 end
329 return {'wait', revision(), kind, tostring(value)}
330end
331
332if operation == 'initialize' then
333 if #ARGV ~= 2 then
334 return failure('INVALID_ARGUMENT', 'Invalid initialize arguments')
335 end
336 ensure_meta()
337 return {'ok'}
338end
339
340if operation == 'put' then
341 if metadata_closed then
342 return failure('FAILED_PRECONDITION', 'Chunk store is closed for writes')
343 end
344 local count = canonical_uint(ARGV[3], max_count_value, '4294967296')
345 if count == nil then
346 return failure('INVALID_ARGUMENT', 'Invalid fragment count')
347 end
348 if #ARGV ~= 3 + count * 5 then
349 return failure('INVALID_ARGUMENT', 'Fragment count does not match arguments')
350 end
351 if count == 0 then
352 return {'ok'}
353 end
354 if metadata_put_count + count > max_count_value then
355 return failure('RESOURCE_EXHAUSTED',
356 'Maximum chunk-store cardinality exceeded')
357 end
358 local explicit = ARGV[4] ~= ''
359 local put_count = metadata_put_count
360 local candidate = put_count
361 local assigned = {}
362 local seen = {}
363 local batch_final = nil
364 local saw_final = false
365 local pending_max = metadata_max_seq
366
367 for index = 1, count do
368 local base = 4 + (index - 1) * 5
369 local supplied_seq = ARGV[base]
370 local is_final = ARGV[base + 1]
371 local storage = ARGV[base + 2]
372 local payload = ARGV[base + 3]
373 local tombstone = ARGV[base + 4]
374 if (supplied_seq ~= '') ~= explicit then
375 return failure('INVALID_ARGUMENT',
376 'Sequence numbers must be set on every fragment or none')
377 end
378 local seq = nil
379 if explicit then
380 seq = canonical_uint(supplied_seq, max_seq_value, '4294967295')
381 if seq == nil then
382 return failure('INVALID_ARGUMENT', 'Invalid explicit sequence number')
383 end
384 else
385 while candidate <= max_seq_value and
386 redis.call('HEXISTS', seq_key, tostring(candidate)) == 1 do
387 candidate = candidate + 1
388 end
389 if candidate > max_seq_value then
390 return failure('RESOURCE_EXHAUSTED',
391 'Maximum implicit sequence number exceeded')
392 end
393 seq = candidate
394 candidate = candidate + 1
395 end
396 local seq_text = tostring(seq)
397 if seen[seq_text] then
398 return failure('INVALID_ARGUMENT',
399 'A sequence occurs more than once in the batch')
400 end
401 seen[seq_text] = true
402 if redis.call('HEXISTS', seq_key, seq_text) == 1 then
403 return failure('ALREADY_EXISTS',
404 'A fragment with seq ' .. seq_text .. ' already exists')
405 end
406 if redis.call('HEXISTS', blobs_key, seq_text) == 1 then
407 return failure('DATA_LOSS',
408 'An unindexed blob collides with seq ' .. seq_text)
409 end
410 local arrival = tostring(put_count + index - 1)
411 if redis.call('HEXISTS', arrival_key, arrival) == 1 then
412 return failure('DATA_LOSS',
413 'Arrival index ' .. arrival .. ' already exists')
414 end
415 if storage ~= 'inline' and storage ~= 'redis' then
416 return failure('INVALID_ARGUMENT', 'Invalid chunk storage kind')
417 end
418 if payload == nil then
419 return failure('INVALID_ARGUMENT', 'Chunk payload is missing')
420 end
421 if tombstone == nil then
422 return failure('INVALID_ARGUMENT', 'Chunk tombstone payload is missing')
423 end
424 if is_final ~= '0' and is_final ~= '1' then
425 return failure('INVALID_ARGUMENT', 'Invalid final-fragment marker')
426 end
427 if is_final == '1' then
428 if saw_final then
429 return failure('INVALID_ARGUMENT',
430 'More than one fragment in the batch is marked final')
431 end
432 if not explicit and index ~= count then
433 return failure('INVALID_ARGUMENT',
434 'The final implicit fragment must be last')
435 end
436 saw_final = true
437 batch_final = seq
438 end
439 assigned[index] = seq
440 if not pending_max or seq > pending_max then
441 pending_max = seq
442 end
443 end
444
445 local existing_final = metadata_final_seq
446 if batch_final and existing_final and batch_final ~= existing_final then
447 return failure('FAILED_PRECONDITION',
448 'The chunk store already has a different final sequence')
449 end
450 local pending_final = batch_final or existing_final
451 if pending_final then
452 local existing_max = metadata_max_seq
453 if existing_max and existing_max > pending_final then
454 return failure('INVALID_ARGUMENT',
455 'An existing fragment exceeds the proposed final sequence')
456 end
457 for index = 1, count do
458 if assigned[index] > pending_final then
459 return failure('INVALID_ARGUMENT',
460 'A fragment sequence exceeds the final sequence')
461 end
462 end
463 end
464
465 local revision_error = can_publish_change()
466 if revision_error then
467 return revision_error
468 end
469 if not can_append_stream_entry() then
470 return failure('RESOURCE_EXHAUSTED', 'Redis Stream ID space is exhausted')
471 end
472 ensure_meta()
473 local response = {'ok'}
474 for index = 1, count do
475 local base = 4 + (index - 1) * 5
476 local storage = ARGV[base + 2]
477 local payload = ARGV[base + 3]
478 local tombstone = ARGV[base + 4]
479 local seq_text = tostring(assigned[index])
480 local arrival = tostring(put_count + index - 1)
481 local stream_id = nil
482 if storage == 'redis' then
483 redis.call('HSET', blobs_key, seq_text, payload)
484 stream_id = redis.call('XADD', stream_key, '*',
485 'v', '1', 'kind', 'chunk', 'seq', seq_text, 'arrival', arrival,
486 'storage', 'redis', 'ref', seq_text, 'tombstone', tombstone)
487 else
488 stream_id = redis.call('XADD', stream_key, '*',
489 'v', '1', 'kind', 'chunk', 'seq', seq_text, 'arrival', arrival,
490 'storage', 'inline', 'payload', payload, 'tombstone', tombstone)
491 end
492 redis.call('HSET', seq_key, seq_text, stream_id)
493 redis.call('HSET', arrival_key, arrival, seq_text)
494 response[#response + 1] = seq_text
495 end
496 redis.call('HINCRBY', meta_key, 'size', count)
497 redis.call('HINCRBY', meta_key, 'put_count', count)
498 redis.call('HSET', meta_key, 'max_seq', tostring(pending_max))
499 if batch_final then
500 redis.call('HSET', meta_key, 'final_seq', tostring(batch_final))
501 end
502 publish_change()
503 return response
504end
505
506if operation == 'lookup' then
507 if #ARGV ~= 4 then
508 return failure('INVALID_ARGUMENT', 'Invalid lookup arguments')
509 end
510 local kind = ARGV[3]
511 local value = ARGV[4]
512 local seq_text = value
513 if kind == 'arrival' then
514 if not is_canonical_decimal(value, '18446744073709551615') then
515 return failure('INVALID_ARGUMENT', 'Invalid arrival-order lookup')
516 end
517 seq_text = redis.call('HGET', arrival_key, value)
518 if not seq_text then
519 return missing_result(kind, value)
520 end
521 elseif kind ~= 'sequence' then
522 return failure('INVALID_ARGUMENT', 'Invalid Redis chunk lookup kind')
523 end
524 local seq = canonical_uint(seq_text, max_seq_value, '4294967295')
525 if seq == nil then
526 if kind == 'arrival' then
527 return failure('DATA_LOSS', 'Arrival index contains an invalid sequence')
528 end
529 return failure('INVALID_ARGUMENT', 'Invalid lookup sequence')
530 end
531 local state, value_or_error, _, _, storage = load_chunk(seq)
532 if state == 'missing' then
533 return missing_result(kind, value)
534 end
535 if state == 'data_loss' then
536 return failure('DATA_LOSS', value_or_error)
537 end
538 if state == 's3' then
539 return {'item', seq_text, 's3', value_or_error, final_seq()}
540 end
541 return {'item', seq_text, storage, value_or_error, final_seq()}
542end
543
544if operation == 'next' then
545 if #ARGV ~= 3 then
546 return failure('INVALID_ARGUMENT', 'Invalid next arguments')
547 end
548 local limit = canonical_uint(ARGV[3], 1024, '1024')
549 if limit == nil or limit == 0 then
550 return failure('INVALID_ARGUMENT', 'limit must be positive')
551 end
552 if not meta_exists then
553 return {'next', 'wait', '', '0'}
554 end
555 local cursor = metadata_next_cursor
556 local final_text = final_seq()
557 local final = metadata_final_seq
558 local items = {}
559 local item_count = 0
560 local disposition = nil
561 local detail = ''
562
563 while true do
564 -- LocalChunkStore treats exhausting the uint32_t sequence namespace as an
565 -- end sentinel even when no final fragment or close status was recorded.
566 if cursor > max_seq_value then
567 disposition = 'end'
568 break
569 end
570 if final and cursor > final then
571 if redis.call('HGET', meta_key, 'closed') == '1' then
572 disposition = 'closed'
573 detail = redis.call('HGET', meta_key, 'status')
574 else
575 disposition = 'end'
576 end
577 break
578 end
579 if item_count == limit then
580 disposition = 'ready'
581 break
582 end
583 local state, value_or_error, _, _, storage = load_chunk(cursor)
584 if state == 'missing' then
585 if redis.call('HGET', meta_key, 'closed') == '1' then
586 disposition = 'closed'
587 detail = redis.call('HGET', meta_key, 'status')
588 else
589 disposition = 'wait'
590 end
591 break
592 end
593 if state == 'data_loss' then
594 disposition = 'data_loss'
595 detail = value_or_error
596 break
597 end
598 if state == 's3' then
599 disposition = 's3'
600 detail = value_or_error
601 break
602 end
603 item_count = item_count + 1
604 items[#items + 1] = tostring(cursor)
605 items[#items + 1] = storage
606 items[#items + 1] = value_or_error
607 items[#items + 1] = final_text
608 cursor = cursor + 1
609 end
610
611 if item_count > 0 and disposition ~= 'data_loss' and disposition ~= 's3' then
612 redis.call('HSET', meta_key, 'next_cursor', tostring(cursor))
613 end
614 local response = {'next', disposition, detail, tostring(item_count)}
615 for index = 1, #items do
616 response[#response + 1] = items[index]
617 end
618 return response
619end
620
621if operation == 'clear' then
622 if #ARGV ~= 3 then
623 return failure('INVALID_ARGUMENT', 'Invalid clear arguments')
624 end
625 local seq_text = ARGV[3]
626 local seq = canonical_uint(seq_text, max_seq_value, '4294967295')
627 if seq == nil then
628 return failure('INVALID_ARGUMENT', 'Invalid sequence to clear')
629 end
630 local state, value_or_error, stream_id, arrival, storage = load_chunk(seq)
631 if state == 'missing' then
632 return failure('NOT_FOUND', 'No fragment with seq ' .. seq_text .. ' exists')
633 end
634 if state == 'data_loss' then
635 return failure('DATA_LOSS', value_or_error)
636 end
637 if state == 's3' then
638 return failure('UNIMPLEMENTED',
639 'Clearing S3-backed chunks is not implemented')
640 end
641 if storage == 'tombstone' then
642 return {'item', seq_text, 'tombstone', value_or_error, final_seq()}
643 end
644 local entry = redis.call('XRANGE', stream_key, stream_id, stream_id,
645 'COUNT', 1)
646 local fields = field_map(entry)
647 if not fields or fields['tombstone'] == nil then
648 return failure('DATA_LOSS',
649 'Stream entry has no prepared tombstone payload')
650 end
651 local revision_error = can_publish_change()
652 if revision_error then
653 return revision_error
654 end
655 if not can_append_stream_entry() then
656 return failure('RESOURCE_EXHAUSTED', 'Redis Stream ID space is exhausted')
657 end
658 local tombstone = fields['tombstone']
659 -- Append first: XADD is the only remaining command that can reject valid
660 -- arguments because of stream-ID state. No old payload has been removed if
661 -- that happens.
662 local tombstone_id = redis.call('XADD', stream_key, '*',
663 'v', '1', 'kind', 'chunk', 'seq', seq_text, 'arrival', arrival,
664 'storage', 'tombstone', 'payload', tombstone,
665 'tombstone', tombstone)
666 if storage == 'redis' then
667 redis.call('HDEL', blobs_key, seq_text)
668 end
669 redis.call('XDEL', stream_key, stream_id)
670 redis.call('HSET', seq_key, seq_text, tombstone_id)
671 publish_change()
672 return {'item', seq_text, storage, value_or_error, final_seq()}
673end
674
675if operation == 'arrival_seq' then
676 if #ARGV ~= 3 then
677 return failure('INVALID_ARGUMENT', 'Invalid arrival lookup arguments')
678 end
679 if not is_canonical_decimal(ARGV[3], '18446744073709551615') then
680 return failure('INVALID_ARGUMENT', 'Invalid arrival order')
681 end
682 local seq = redis.call('HGET', arrival_key, ARGV[3])
683 if not seq then
684 return failure('NOT_FOUND',
685 'No fragment has arrival order ' .. ARGV[3])
686 end
687 if canonical_uint(seq, max_seq_value, '4294967295') == nil or
688 redis.call('HEXISTS', seq_key, seq) ~= 1 then
689 return failure('DATA_LOSS',
690 'Arrival index references an invalid or missing sequence')
691 end
692 return {'value', seq}
693end
694
695if operation == 'final' then
696 if #ARGV ~= 2 then
697 return failure('INVALID_ARGUMENT', 'Invalid final-sequence arguments')
698 end
699 return {'optional', final_seq()}
700end
701
702if operation == 'size' then
703 if #ARGV ~= 2 then
704 return failure('INVALID_ARGUMENT', 'Invalid size arguments')
705 end
706 if not meta_exists then
707 return {'value', '0'}
708 end
709 return {'value', redis.call('HGET', meta_key, 'size')}
710end
711
712if operation == 'close' then
713 if #ARGV ~= 4 then
714 return failure('INVALID_ARGUMENT', 'Invalid close arguments')
715 end
716 local requested_status = ARGV[3]
717 local return_existing = ARGV[4]
718 if requested_status == nil or
719 (return_existing ~= '0' and return_existing ~= '1') then
720 return failure('INVALID_ARGUMENT', 'Invalid chunk-store close arguments')
721 end
722 if metadata_closed then
723 if return_existing == '1' then
724 return {'status', redis.call('HGET', meta_key, 'status')}
725 end
726 return failure('FAILED_PRECONDITION',
727 'Chunk store is already closed for writes')
728 end
729 local revision_error = can_publish_change()
730 if revision_error then
731 return revision_error
732 end
733 if not can_append_stream_entry() then
734 return failure('RESOURCE_EXHAUSTED', 'Redis Stream ID space is exhausted')
735 end
736 local changed = metadata_revision + 1
737 -- As in ClearData, append the stream transition before making any other
738 -- mutation so an invalid stream-ID state cannot leave a half-closed store.
739 redis.call('XADD', stream_key, '*', 'v', '1', 'kind', 'control',
740 'event', 'close', 'revision', tostring(changed))
741 ensure_meta()
742 redis.call('HSET', meta_key, 'closed', '1', 'status', requested_status)
743 redis.call('HINCRBY', meta_key, 'revision', 1)
744 redis.call('PUBLISH', events_channel, tostring(changed))
745 return {'status', requested_status}
746end
747
748if operation == 'metadata' then
749 if #ARGV ~= 2 then
750 return failure('INVALID_ARGUMENT', 'Invalid metadata arguments')
751 end
752 if not meta_exists then
753 return {'metadata', node_id, '0', '', '', '0', '0', '0', '', '0'}
754 end
755 local closed = redis.call('HGET', meta_key, 'closed')
756 return {'metadata',
757 redis.call('HGET', meta_key, 'id'),
758 closed,
759 closed == '1' and redis.call('HGET', meta_key, 'status') or '',
760 redis.call('HGET', meta_key, 'final_seq') or '',
761 redis.call('HGET', meta_key, 'size'),
762 redis.call('HGET', meta_key, 'put_count'),
763 redis.call('HGET', meta_key, 'next_cursor'),
764 redis.call('HGET', meta_key, 'max_seq') or '',
765 redis.call('HGET', meta_key, 'revision')}
766end
767
768return failure('INVALID_ARGUMENT', 'Unknown chunk store operation')
769)lua";
770
771} // namespace a11::stores::internal
772
773#endif // A11_STORES_REDIS_CHUNK_STORE_SCRIPT_H_
Definition redis_chunk_store_script.h:22
constexpr std::string_view kRedisChunkStoreScript
Definition redis_chunk_store_script.h:27