feat: preserve stereo channels and require reviewed speaker roles
This commit is contained in:
@@ -7,19 +7,34 @@ namespace app\common\service\followupaudio;
|
||||
/** Versioned encrypted state. Only completed responses, never in-flight intents, are recoverable. */
|
||||
final class FollowupAudioPipelineCheckpoint
|
||||
{
|
||||
public static function initial(array $task, string $fingerprint, int $chunkMs): array
|
||||
public const CHANNEL_POLICY = 'separate-fixed-windows-v1';
|
||||
|
||||
public static function initial(array $task, string $fingerprint, int $chunkMs, int $sourceChannels = 1): array
|
||||
{
|
||||
return self::plan($task, $fingerprint, $chunkMs, $sourceChannels, 2);
|
||||
}
|
||||
|
||||
/** V1 reconstruction exists only for reading immutable historical checkpoints. */
|
||||
private static function plan(array $task, string $fingerprint, int $chunkMs, int $sourceChannels, int $version): array
|
||||
{
|
||||
$duration = (int) ceil((float) $task['duration_seconds'] * 1000);
|
||||
if ($duration < 1 || $duration > 3600000 || $chunkMs < 1000 || $chunkMs > 120000) { self::invalid(); }
|
||||
FollowupAudioPipelineMedia::assertChunkPlan((float) $task['duration_seconds'], $chunkMs);
|
||||
FollowupAudioPipelineMedia::assertChunkPlan((float) $task['duration_seconds'], $chunkMs, $sourceChannels);
|
||||
$chunks = [];
|
||||
for ($start = 0; $start < $duration; $start += $chunkMs) {
|
||||
$end = min($duration, $start + $chunkMs);
|
||||
$chunks[] = ['id' => 'asr-' . hash('sha256', $task['sha256'] . ':' . $fingerprint . ':' . $start . ':' . $end),
|
||||
'start_ms' => $start, 'end_ms' => $end, 'state' => 'pending'];
|
||||
for ($channel = 0; $channel < $sourceChannels; $channel++) {
|
||||
$identity = $task['sha256'] . ':' . $fingerprint . ':' . $start . ':' . $end;
|
||||
if ($version === 2) { $identity .= ':' . self::CHANNEL_POLICY . ':' . $sourceChannels . ':' . $channel; }
|
||||
$entry = ['id' => 'asr-' . hash('sha256', $identity), 'start_ms' => $start, 'end_ms' => $end, 'state' => 'pending'];
|
||||
if ($sourceChannels === 2) { $entry['channel'] = $channel; }
|
||||
$chunks[] = $entry;
|
||||
}
|
||||
}
|
||||
return ['schema_version' => 1, 'revision' => 0, 'source_sha256' => $task['sha256'], 'fingerprint' => $fingerprint,
|
||||
$state = ['schema_version' => $version, 'revision' => 0, 'source_sha256' => $task['sha256'], 'fingerprint' => $fingerprint,
|
||||
'duration_ms' => $duration, 'chunk_ms' => $chunkMs, 'chunks' => $chunks, 'extraction' => ['state' => 'pending']];
|
||||
if ($version === 2) { $state += ['source_channels' => $sourceChannels, 'channel_policy' => self::CHANNEL_POLICY, 'silence_policy' => FollowupAudioPipelineMedia::SILENCE_POLICY]; }
|
||||
return $state;
|
||||
}
|
||||
|
||||
public static function validate(array $state, array $task): void
|
||||
@@ -30,9 +45,11 @@ final class FollowupAudioPipelineCheckpoint
|
||||
|| ($state['source_sha256'] ?? '') !== ($task['sha256'] ?? '')
|
||||
|| !is_int($state['revision'] ?? null) || $state['revision'] < 1 || $state['revision'] > 15000
|
||||
|| !is_int($state['chunk_ms'] ?? null)) { self::invalid(); }
|
||||
$expected = self::initial($task, $state['fingerprint'], $state['chunk_ms']);
|
||||
foreach (['schema_version', 'source_sha256', 'fingerprint', 'duration_ms', 'chunk_ms'] as $field) {
|
||||
if (($state[$field] ?? null) !== $expected[$field]) { self::invalid(); }
|
||||
$version = $state['schema_version'] ?? 0;
|
||||
if (!in_array($version, [1, 2], true) || ($version === 2 && !is_int($state['source_channels'] ?? null))) { self::invalid(); }
|
||||
$expected = self::plan($task, $state['fingerprint'], $state['chunk_ms'], $version === 2 ? $state['source_channels'] : 1, $version);
|
||||
foreach (['schema_version', 'source_sha256', 'fingerprint', 'duration_ms', 'chunk_ms', 'source_channels', 'channel_policy', 'silence_policy'] as $field) {
|
||||
if (($state[$field] ?? null) !== ($expected[$field] ?? null)) { self::invalid(); }
|
||||
}
|
||||
if (array_diff(array_keys($state), array_keys($expected)) || !is_array($state['chunks'] ?? null)
|
||||
|| !array_is_list($state['chunks']) || count($state['chunks']) !== count($expected['chunks'])) { self::invalid(); }
|
||||
@@ -40,10 +57,11 @@ final class FollowupAudioPipelineCheckpoint
|
||||
$unfinished = false;
|
||||
foreach ($state['chunks'] as $index => $chunk) {
|
||||
if (!is_array($chunk)) { self::invalid(); }
|
||||
foreach (['id', 'start_ms', 'end_ms'] as $field) { if (($chunk[$field] ?? null) !== $expected['chunks'][$index][$field]) { self::invalid(); } }
|
||||
foreach (['id', 'start_ms', 'end_ms', 'channel'] as $field) { if (($chunk[$field] ?? null) !== ($expected['chunks'][$index][$field] ?? null)) { self::invalid(); } }
|
||||
if ($version === 1 && $chunk['state'] === 'local_silence') { self::invalid(); }
|
||||
self::entry($chunk, false);
|
||||
if ($unfinished && $chunk['state'] !== 'pending') { self::invalid(); }
|
||||
if ($chunk['state'] !== 'complete') { $unfinished = true; }
|
||||
if (!in_array($chunk['state'], ['complete', 'local_silence'], true)) { $unfinished = true; }
|
||||
$textBytes += strlen($chunk['text'] ?? '');
|
||||
}
|
||||
if ($textBytes > 2000000 || !is_array($state['extraction'] ?? null)) { self::invalid(); }
|
||||
@@ -53,8 +71,18 @@ final class FollowupAudioPipelineCheckpoint
|
||||
|
||||
private static function entry(array $entry, bool $extraction): void
|
||||
{
|
||||
$allowed = $extraction ? ['state', 'request_id', 'answer'] : ['id', 'start_ms', 'end_ms', 'state', 'request_id', 'text'];
|
||||
if (array_diff(array_keys($entry), $allowed) || !in_array($entry['state'] ?? '', ['pending', 'intent', 'complete', 'rejected'], true)) { self::invalid(); }
|
||||
$allowed = $extraction ? ['state', 'request_id', 'answer'] : ['id', 'start_ms', 'end_ms', 'channel', 'state', 'request_id', 'text', 'source', 'pcm_sha256', 'pcm_bytes'];
|
||||
if (array_diff(array_keys($entry), $allowed) || !in_array($entry['state'] ?? '', ['pending', 'intent', 'complete', 'rejected', 'local_silence'], true)) { self::invalid(); }
|
||||
if (array_key_exists('channel', $entry) && (!is_int($entry['channel']) || !in_array($entry['channel'], [0, 1], true))) { self::invalid(); }
|
||||
if ($entry['state'] === 'local_silence') {
|
||||
if ($extraction || array_key_exists('request_id', $entry) || ($entry['text'] ?? null) !== '' || ($entry['source'] ?? '') !== 'local_silence'
|
||||
|| !is_int($entry['pcm_bytes'] ?? null) || $entry['pcm_bytes'] < 2 || $entry['pcm_bytes'] > 3842560 || $entry['pcm_bytes'] % 2
|
||||
|| !is_string($entry['pcm_sha256'] ?? null) || !hash_equals(self::zeroHash($entry['pcm_bytes']), $entry['pcm_sha256'])) { self::invalid(); }
|
||||
$expected = ($entry['end_ms'] - $entry['start_ms']) * 32;
|
||||
if (abs($entry['pcm_bytes'] - $expected) > 2560) { self::invalid(); }
|
||||
return;
|
||||
}
|
||||
if (array_key_exists('source', $entry) || array_key_exists('pcm_sha256', $entry) || array_key_exists('pcm_bytes', $entry)) { self::invalid(); }
|
||||
if ($entry['state'] === 'pending') {
|
||||
if (array_key_exists('request_id', $entry) || array_key_exists($extraction ? 'answer' : 'text', $entry)) { self::invalid(); }
|
||||
} elseif (!is_string($entry['request_id'] ?? null) || !preg_match('/^fa-[a-f0-9]{32}$/D', $entry['request_id'])) { self::invalid(); }
|
||||
@@ -69,7 +97,7 @@ final class FollowupAudioPipelineCheckpoint
|
||||
{
|
||||
self::validate($next, $task);
|
||||
if ($previous === []) {
|
||||
$initial = self::initial($task, $next['fingerprint'], $next['chunk_ms']);
|
||||
$initial = self::initial($task, $next['fingerprint'], $next['chunk_ms'], $next['source_channels'] ?? 1);
|
||||
$initial['revision'] = 1;
|
||||
if ($expectedRevision !== 0 || FollowupAudioPolicy::canonical($next) !== FollowupAudioPolicy::canonical($initial)) { self::invalid(); }
|
||||
return;
|
||||
@@ -78,18 +106,19 @@ final class FollowupAudioPipelineCheckpoint
|
||||
if ($expectedRevision !== $previous['revision'] || $next['revision'] !== $previous['revision'] + 1) {
|
||||
throw new FollowupAudioException('PIPELINE_CHECKPOINT_CONFLICT');
|
||||
}
|
||||
foreach (['schema_version', 'source_sha256', 'fingerprint', 'duration_ms', 'chunk_ms'] as $key) {
|
||||
if ($next[$key] !== $previous[$key]) { self::invalid(); }
|
||||
foreach (['schema_version', 'source_sha256', 'fingerprint', 'duration_ms', 'chunk_ms', 'source_channels', 'channel_policy', 'silence_policy'] as $key) {
|
||||
if (($next[$key] ?? null) !== ($previous[$key] ?? null)) { self::invalid(); }
|
||||
}
|
||||
$changed = 0;
|
||||
foreach (array_merge($previous['chunks'], [$previous['extraction']]) as $index => $before) {
|
||||
$after = $index < count($next['chunks']) ? $next['chunks'][$index] : $next['extraction'];
|
||||
if (FollowupAudioPolicy::canonical($before) === FollowupAudioPolicy::canonical($after)) { continue; }
|
||||
$changed++;
|
||||
foreach (['id', 'start_ms', 'end_ms'] as $key) { if (($before[$key] ?? null) !== ($after[$key] ?? null)) { self::invalid(); } }
|
||||
foreach (['id', 'start_ms', 'end_ms', 'channel'] as $key) { if (($before[$key] ?? null) !== ($after[$key] ?? null)) { self::invalid(); } }
|
||||
$from = $before['state']; $to = $after['state'];
|
||||
if (!((in_array($from, ['pending', 'rejected'], true) && $to === 'intent')
|
||||
|| ($from === 'intent' && in_array($to, ['complete', 'rejected'], true)))) { self::invalid(); }
|
||||
|| ($from === 'intent' && in_array($to, ['complete', 'rejected'], true))
|
||||
|| ($from === 'pending' && $to === 'local_silence' && $next['schema_version'] === 2))) { self::invalid(); }
|
||||
if ($from === 'intent' && $before['request_id'] !== $after['request_id']) { self::invalid(); }
|
||||
}
|
||||
if ($changed !== 1) { self::invalid(); }
|
||||
@@ -105,10 +134,16 @@ final class FollowupAudioPipelineCheckpoint
|
||||
{
|
||||
$result = [];
|
||||
foreach ($state['chunks'] as $chunk) {
|
||||
if ($chunk['state'] === 'complete') { $result[] = array_intersect_key($chunk, array_flip(['id', 'start_ms', 'end_ms', 'text'])); }
|
||||
if (in_array($chunk['state'], ['complete', 'local_silence'], true)) { $result[] = array_intersect_key($chunk, array_flip(['id', 'start_ms', 'end_ms', 'text', 'channel', 'source'])); }
|
||||
}
|
||||
return $result;
|
||||
}
|
||||
|
||||
private static function zeroHash(int $bytes): string
|
||||
{
|
||||
static $cache = [];
|
||||
return $cache[$bytes] ?? ($cache[$bytes] = hash('sha256', str_repeat("\0", $bytes)));
|
||||
}
|
||||
|
||||
private static function invalid(): void { throw new FollowupAudioException('PIPELINE_CHECKPOINT_INVALID', true); }
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user