17#ifndef A11_STORES_REDIS_CHUNK_STORE_SCRIPT_H_
18#define A11_STORES_REDIS_CHUNK_STORE_SCRIPT_H_
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]
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
44if type(operation) ~= 'string' or type(node_id) ~= 'string' or node_id == '' then
45 return {'error', 'INVALID_ARGUMENT', 'Invalid chunk-store script envelope'}
48local function failure(code, message)
49 return {'error', code, message}
52local function key_type(key)
53 local reply = redis.call('TYPE', key)
54 if type(reply) == 'table' then
60local function check_type(key, expected)
61 local actual = key_type(key)
62 return actual == 'none' or actual == expected
65local function is_canonical_decimal(value, maximum)
66 if type(value) ~= 'string' or value == '' or
67 string.find(value, '[^0-9]') then
70 if #value > 1 and string.sub(value, 1, 1) == '0' then
73 if #value ~= #maximum then
74 return #value < #maximum
76 return value <= maximum
79local function canonical_uint(value, maximum, maximum_text)
80 if not is_canonical_decimal(value, maximum_text) then
83 local parsed = tonumber(value)
84 if not parsed or parsed < 0 or parsed > maximum or
85 math.floor(parsed) ~= parsed then
91if not check_type(meta_key, 'hash') then
92 return failure('DATA_LOSS', 'Chunk store metadata key is not a hash')
94if not check_type(stream_key, 'stream') then
95 return failure('DATA_LOSS', 'Chunk store data key is not a stream')
97if not check_type(seq_key, 'hash') then
98 return failure('DATA_LOSS', 'Chunk store sequence index is not a hash')
100if not check_type(arrival_key, 'hash') then
101 return failure('DATA_LOSS', 'Chunk store arrival index is not a hash')
103if not check_type(blobs_key, 'hash') then
104 return failure('DATA_LOSS', 'Chunk store blob key is not a hash')
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')
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')
137 metadata_final_seq = canonical_uint(
138 final_seq, max_seq_value, '4294967295')
141 metadata_max_seq = canonical_uint(max_seq, max_seq_value, '4294967295')
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')
151 if stored_id ~= node_id then
152 return failure('DATA_LOSS', 'Chunk store key belongs to a different node')
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')
159 if not metadata_closed and stored_status then
160 return failure('DATA_LOSS', 'Open chunk store has a terminal status')
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')
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')
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')
176 if redis.call('HLEN', blobs_key) > metadata_size then
177 return failure('DATA_LOSS', 'Chunk store contains unindexed blobs')
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')
185local function ensure_meta()
186 if not meta_exists then
187 redis.call('HSET', meta_key,
199local function revision()
200 if not meta_exists then
203 return redis.call('HGET', meta_key, 'revision') or '0'
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))
214local function can_publish_change()
215 if metadata_revision >= max_revision_value then
216 return failure('RESOURCE_EXHAUSTED',
217 'Chunk store mutation revision exhausted')
222local function field_map(entry)
226 local fields = entry[1][2]
228 for index = 1, #fields, 2 do
229 result[fields[index]] = fields[index + 1]
234local function can_append_stream_entry()
235 if redis.call('EXISTS', stream_key) == 0 then
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]
246 if not stream_id then
249 local separator = string.find(stream_id, '-', 1, true)
250 if not separator then
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'
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', '', '', '', ''
267 local entry = redis.call('XRANGE', stream_key, stream_id, stream_id, 'COUNT', 1)
268 local fields = field_map(entry)
270 return 'data_loss', 'Sequence index references a missing stream entry',
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, '', ''
277 local arrival = canonical_uint(
278 fields['arrival'], max_seq_value, '4294967295')
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'], ''
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
290 return 'item', fields['payload'], stream_id, fields['arrival'], storage
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
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
302 return 'item', payload, stream_id, fields['arrival'], storage
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
309 return 'item', fields['payload'], stream_id, fields['arrival'], storage
311 if storage == 's3' then
312 return 's3', fields['ref'] or '', stream_id, fields['arrival'], storage
314 return 'data_loss', 'Stream entry has an unknown storage kind', stream_id,
315 fields['arrival'], storage or ''
318local function final_seq()
319 if not meta_exists then
322 return redis.call('HGET', meta_key, 'final_seq') or ''
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')}
329 return {'wait', revision(), kind, tostring(value)}
332if operation == 'initialize' then
334 return failure('INVALID_ARGUMENT', 'Invalid initialize arguments')
340if operation == 'put' then
341 if metadata_closed then
342 return failure('FAILED_PRECONDITION', 'Chunk store is closed for writes')
344 local count = canonical_uint(ARGV[3], max_count_value, '4294967296')
346 return failure('INVALID_ARGUMENT', 'Invalid fragment count')
348 if #ARGV ~= 3 + count * 5 then
349 return failure('INVALID_ARGUMENT', 'Fragment count does not match arguments')
354 if metadata_put_count + count > max_count_value then
355 return failure('RESOURCE_EXHAUSTED',
356 'Maximum chunk-store cardinality exceeded')
358 local explicit = ARGV[4] ~= ''
359 local put_count = metadata_put_count
360 local candidate = put_count
363 local batch_final = nil
364 local saw_final = false
365 local pending_max = metadata_max_seq
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')
380 seq = canonical_uint(supplied_seq, max_seq_value, '4294967295')
382 return failure('INVALID_ARGUMENT', 'Invalid explicit sequence number')
385 while candidate <= max_seq_value and
386 redis.call('HEXISTS', seq_key, tostring(candidate)) == 1 do
387 candidate = candidate + 1
389 if candidate > max_seq_value then
390 return failure('RESOURCE_EXHAUSTED',
391 'Maximum implicit sequence number exceeded')
394 candidate = candidate + 1
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')
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')
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)
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')
415 if storage ~= 'inline' and storage ~= 'redis' then
416 return failure('INVALID_ARGUMENT', 'Invalid chunk storage kind')
418 if payload == nil then
419 return failure('INVALID_ARGUMENT', 'Chunk payload is missing')
421 if tombstone == nil then
422 return failure('INVALID_ARGUMENT', 'Chunk tombstone payload is missing')
424 if is_final ~= '0' and is_final ~= '1' then
425 return failure('INVALID_ARGUMENT', 'Invalid final-fragment marker')
427 if is_final == '1' then
429 return failure('INVALID_ARGUMENT',
430 'More than one fragment in the batch is marked final')
432 if not explicit and index ~= count then
433 return failure('INVALID_ARGUMENT',
434 'The final implicit fragment must be last')
439 assigned[index] = seq
440 if not pending_max or seq > pending_max then
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')
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')
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')
465 local revision_error = can_publish_change()
466 if revision_error then
467 return revision_error
469 if not can_append_stream_entry() then
470 return failure('RESOURCE_EXHAUSTED', 'Redis Stream ID space is exhausted')
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)
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)
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
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))
500 redis.call('HSET', meta_key, 'final_seq', tostring(batch_final))
506if operation == 'lookup' then
508 return failure('INVALID_ARGUMENT', 'Invalid lookup arguments')
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')
517 seq_text = redis.call('HGET', arrival_key, value)
519 return missing_result(kind, value)
521 elseif kind ~= 'sequence' then
522 return failure('INVALID_ARGUMENT', 'Invalid Redis chunk lookup kind')
524 local seq = canonical_uint(seq_text, max_seq_value, '4294967295')
526 if kind == 'arrival' then
527 return failure('DATA_LOSS', 'Arrival index contains an invalid sequence')
529 return failure('INVALID_ARGUMENT', 'Invalid lookup sequence')
531 local state, value_or_error, _, _, storage = load_chunk(seq)
532 if state == 'missing' then
533 return missing_result(kind, value)
535 if state == 'data_loss' then
536 return failure('DATA_LOSS', value_or_error)
538 if state == 's3' then
539 return {'item', seq_text, 's3', value_or_error, final_seq()}
541 return {'item', seq_text, storage, value_or_error, final_seq()}
544if operation == 'next' then
546 return failure('INVALID_ARGUMENT', 'Invalid next arguments')
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')
552 if not meta_exists then
553 return {'next', 'wait', '', '0'}
555 local cursor = metadata_next_cursor
556 local final_text = final_seq()
557 local final = metadata_final_seq
560 local disposition = nil
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
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')
579 if item_count == limit then
580 disposition = 'ready'
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')
593 if state == 'data_loss' then
594 disposition = 'data_loss'
595 detail = value_or_error
598 if state == 's3' then
600 detail = value_or_error
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
611 if item_count > 0 and disposition ~= 'data_loss' and disposition ~= 's3' then
612 redis.call('HSET', meta_key, 'next_cursor', tostring(cursor))
614 local response = {'next', disposition, detail, tostring(item_count)}
615 for index = 1, #items do
616 response[#response + 1] = items[index]
621if operation == 'clear' then
623 return failure('INVALID_ARGUMENT', 'Invalid clear arguments')
625 local seq_text = ARGV[3]
626 local seq = canonical_uint(seq_text, max_seq_value, '4294967295')
628 return failure('INVALID_ARGUMENT', 'Invalid sequence to clear')
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')
634 if state == 'data_loss' then
635 return failure('DATA_LOSS', value_or_error)
637 if state == 's3' then
638 return failure('UNIMPLEMENTED',
639 'Clearing S3-backed chunks is not implemented')
641 if storage == 'tombstone' then
642 return {'item', seq_text, 'tombstone', value_or_error, final_seq()}
644 local entry = redis.call('XRANGE', stream_key, stream_id, stream_id,
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')
651 local revision_error = can_publish_change()
652 if revision_error then
653 return revision_error
655 if not can_append_stream_entry() then
656 return failure('RESOURCE_EXHAUSTED', 'Redis Stream ID space is exhausted')
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
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)
669 redis.call('XDEL', stream_key, stream_id)
670 redis.call('HSET', seq_key, seq_text, tombstone_id)
672 return {'item', seq_text, storage, value_or_error, final_seq()}
675if operation == 'arrival_seq' then
677 return failure('INVALID_ARGUMENT', 'Invalid arrival lookup arguments')
679 if not is_canonical_decimal(ARGV[3], '18446744073709551615') then
680 return failure('INVALID_ARGUMENT', 'Invalid arrival order')
682 local seq = redis.call('HGET', arrival_key, ARGV[3])
684 return failure('NOT_FOUND',
685 'No fragment has arrival order ' .. ARGV[3])
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')
692 return {'value', seq}
695if operation == 'final' then
697 return failure('INVALID_ARGUMENT', 'Invalid final-sequence arguments')
699 return {'optional', final_seq()}
702if operation == 'size' then
704 return failure('INVALID_ARGUMENT', 'Invalid size arguments')
706 if not meta_exists then
707 return {'value', '0'}
709 return {'value', redis.call('HGET', meta_key, 'size')}
712if operation == 'close' then
714 return failure('INVALID_ARGUMENT', 'Invalid close arguments')
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')
722 if metadata_closed then
723 if return_existing == '1' then
724 return {'status', redis.call('HGET', meta_key, 'status')}
726 return failure('FAILED_PRECONDITION',
727 'Chunk store is already closed for writes')
729 local revision_error = can_publish_change()
730 if revision_error then
731 return revision_error
733 if not can_append_stream_entry() then
734 return failure('RESOURCE_EXHAUSTED', 'Redis Stream ID space is exhausted')
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))
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}
748if operation == 'metadata' then
750 return failure('INVALID_ARGUMENT', 'Invalid metadata arguments')
752 if not meta_exists then
753 return {'metadata', node_id, '0', '', '', '0', '0', '0', '', '0'}
755 local closed = redis.call('HGET', meta_key, 'closed')
757 redis.call('HGET', meta_key, 'id'),
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')}
768return failure('INVALID_ARGUMENT', 'Unknown chunk store operation')
Definition redis_chunk_store_script.h:22
constexpr std::string_view kRedisChunkStoreScript
Definition redis_chunk_store_script.h:27