Align playback observe API

pull/370/head
zijiren233 4 months ago
parent d3cf4dfdc0
commit 889642a8e5
No known key found for this signature in database
GPG Key ID: 534E082AAA9B39DC

@ -5,7 +5,7 @@ description: 使用 Realtime 订阅播放状态、播放信息、房间设置、
import { Steps, TabItem, Tabs } from '@astrojs/starlight/components';
资源观察用于订阅可缓存查询结果。客户端提交自己的缓存版本,服务端判断是否变化;变化时可以推送完整快照,也可以只通知客户端重新拉取。
资源观察用于订阅房间资源。带版本资源可以携带恢复游标,播放信息始终返回当前播放源可直接交给播放器使用的数据。资源变化时,服务端可以推送完整 payload,也可以只通知客户端重新拉取。
## 支持的资源
@ -33,7 +33,7 @@ import { Steps, TabItem, Tabs } from '@astrojs/starlight/components';
<Steps>
1. 客户端连接房间 WebSocket。
2. 客户端发送 `ClientMessage.observe_resource`,包含 `observe_id`、`delivery_mode` 和目标资源。支持恢复游标的资源还可以带本地缓存版本或事件序列。
2. 客户端发送 `ClientMessage.observe_resource`,包含 `observe_id`、`delivery_mode` 和目标资源。支持恢复游标的资源还可以带事件序列。
3. 服务端加载当前资源,返回 `ServerMessage.resource_observed`。
4. 如果服务端判断资源已变化,随后返回同一 `observe_id` 的 `ServerMessage.resource_changed`。
5. 连接期间,房间事件、缓存失效、Provider credential 变更或播放信息过期会触发重新评估。
@ -44,16 +44,16 @@ import { Steps, TabItem, Tabs } from '@astrojs/starlight/components';
| 字段 | 说明 |
| --- | --- |
| `observe_id` | 对应客户端提交的订阅 ID |
| `version` | 服务端当前版本;播放信息不暴露客户端缓存版本 |
| `changed` | 服务端当前资源是否需要发送给客户端 |
| `event_cursor` | 支持事件回放资源的恢复游标 |
`ResourceChanged` 字段:
| 字段 | 说明 |
| --- | --- |
| `observe_id` | 哪个订阅发生变化 |
| `version` | 新版本;播放信息不使用该字段做客户端缓存恢复 |
| `payload` | 完整快照或 `changed_only`,取决于 `delivery_mode` |
| `event_cursor` | 已回放资源事件对应的游标 |
## 取消订阅
@ -80,9 +80,10 @@ sendClientMessage({
sendClientMessage({
observeResource: {
observeId: 'playback-state',
version: cached.playbackStateVersion ?? '',
deliveryMode: ResourceDeliveryMode.RESOURCE_DELIVERY_MODE_PUSH_SNAPSHOT,
playbackState: {},
playbackState: {
afterEventSequence: cached.playbackStateEventSequence,
},
},
});
```
@ -113,9 +114,9 @@ sendClientMessage({
sendClientMessage({
observeResource: {
observeId: 'room-members:online',
version: cached.membersVersion ?? '',
deliveryMode: ResourceDeliveryMode.RESOURCE_DELIVERY_MODE_NOTIFY_ONLY,
roomMembers: {
afterEventSequence: cached.membersEventSequence,
request: {
page: 1,
pageSize: 50,

@ -22,7 +22,7 @@ Room WebSocket connections carry playback controls, chat, WebRTC signaling, and
2. Fetch current playback state and a playback info.
3. Connect Realtime.
4. Receive play, pause, seek, or media-change events.
5. Refresh the playback info when the media URL expires, Provider credentials change, or resource versions change.
5. Refresh the playback info when the media URL expires, Provider credentials change, or room resources change.
## Direct and Proxy Playback

@ -5,14 +5,14 @@ description: Subscribe to playback state, playbacks, room settings, playlist ite
import { Steps, TabItem, Tabs } from '@astrojs/starlight/components';
Resource observation subscribes to cacheable query results. The client submits its cached version, and the server decides whether the resource changed. When it changes, the server can push a full snapshot or only notify the client to fetch.
Resource observation subscribes to room resources. Versioned resources can carry recovery cursors, while playback info always returns player-ready data for the current source. When a resource changes, the server can push a full payload or only notify the client to fetch.
## Supported Resources
| Resource | `ObserveResource.resource` field | Version source | Typical use |
| --- | --- | --- | --- |
| Playback state | `playback_state` | `RoomPlaybackState.version` | Player UI and synchronized playback status |
| Playback | `playback` | Current value on every observe | Current playback URL, headers, proxy policy, and expiry |
| Playback | `playback` | Current value on every observe | Current playback URL, headers, proxy policy, subtitles, and expiry |
| Room settings | `room_settings` | Room settings version | Chat, playback permissions, and UI controls |
| Playlist items | `playlist_items` | List snapshot version | Root lists, child playlists, provider browsing, and search results |
| Room members | `room_members` | Member list snapshot version | Member pagination, search, role filtering, and status filtering |
@ -33,7 +33,7 @@ Each connection can hold up to 64 observations. Over-limit requests return `Reso
<Steps>
1. The client connects to the room WebSocket.
2. The client sends `ClientMessage.observe_resource` with `observe_id`, `delivery_mode`, and the target resource. Resources that support recovery can also carry a cached version or event sequence.
2. The client sends `ClientMessage.observe_resource` with `observe_id`, `delivery_mode`, and the target resource. Resources that support recovery can also carry an event sequence.
3. The server loads the current resource and sends `ServerMessage.resource_observed`.
4. If the server determines that the resource changed, it sends `ServerMessage.resource_changed` for the same `observe_id`.
5. Later room events, cache invalidation, Provider credential changes, or playback expiry can trigger re-evaluation.
@ -44,16 +44,16 @@ Each connection can hold up to 64 observations. Over-limit requests return `Reso
| Field | Meaning |
| --- | --- |
| `observe_id` | The submitted subscription ID |
| `version` | Current server-side version; playback info does not expose a client cache version |
| `changed` | Whether the server should send the current resource to the client |
| `event_cursor` | Recovery cursor for resources with event replay |
`ResourceChanged` fields:
| Field | Meaning |
| --- | --- |
| `observe_id` | Which subscription changed |
| `version` | New version; playback info does not use it for client cache recovery |
| `payload` | Full snapshot or `changed_only`, depending on `delivery_mode` |
| `event_cursor` | Cursor associated with replayed resource events |
## Unsubscribe
@ -80,9 +80,10 @@ The examples below use TypeScript-style pseudocode. Exact field names depend on
sendClientMessage({
observeResource: {
observeId: 'playback-state',
version: cached.playbackStateVersion ?? '',
deliveryMode: ResourceDeliveryMode.RESOURCE_DELIVERY_MODE_PUSH_SNAPSHOT,
playbackState: {},
playbackState: {
afterEventSequence: cached.playbackStateEventSequence,
},
},
});
```
@ -113,9 +114,9 @@ Playback observation immediately returns the current playable URLs, headers, sub
sendClientMessage({
observeResource: {
observeId: 'room-members:online',
version: cached.membersVersion ?? '',
deliveryMode: ResourceDeliveryMode.RESOURCE_DELIVERY_MODE_NOTIFY_ONLY,
roomMembers: {
afterEventSequence: cached.membersEventSequence,
request: {
page: 1,
pageSize: 50,

@ -16,7 +16,7 @@ Playback issues often sit between the client, SyncTV, Provider, and upstream med
2. SyncTV uses room state, media, user, Provider credentials, and client profile to build `Playback`.
3. The Provider returns one or more `PlaybackInfo` entries with URL, format, headers, subtitles, expiry, and metadata.
4. The client chooses direct URL, proxy URL, transcode variant, HLS, FLV, or subtitle URL.
5. When URL expiry, media switch, credential change, or client capability changes, the client refreshes the snapshot.
5. When URL expiry, media switch, credential change, or client capability changes, the client requests fresh playback info.
</Steps>
## Direct and Proxy Playback
@ -39,7 +39,7 @@ SyncTV proxy uses only headers explicitly provided by the Provider. It does not
| Browser cannot set `Referer` or custom headers | Prefer proxy URL |
| Native client can set headers and access upstream | Direct URL can be used |
| Mobile client lacks codec/container support | Declare capabilities in `PlaybackClientProfile` |
| URL has `expires_at` | Refresh snapshot before expiry |
| URL has `expires_at` | Request fresh playback info before expiry |
| Subtitle has separate headers | Use subtitle-specific headers or proxy subtitle URL |
## Range and Slice Cache
@ -67,7 +67,7 @@ New clients observe playback resources over Realtime:
- `playback_state`: current position, state, and version.
- `playback`: playback URLs, headers, subtitles, and expiry.
On reconnect, send local `version`, `media_id`, `playlist_id`, `target`, and `playback_client_profile`. The server decides whether the cached snapshot is still valid.
On reconnect, fetch `playback_state` and observe `playback` with the current `playback_client_profile`. Playback info is player-ready data for the current source and expires with its URLs.
See [Realtime API](../../develop/realtime-api/).

@ -23,7 +23,7 @@ Users without playback permissions follow the room state. If controls are missin
| --- | --- | --- |
| Playback drift | Refresh or re-enter the room | Room ID, time, whether only you are affected |
| Media switch fails | Retry or refresh playlist | Media ID, error text, HTTP status |
| Playback URL expired | Refresh snapshot or re-enter | Provider name, expiry message, request_id |
| Playback URL expired | Fetch fresh playback info or re-enter | Provider name, expiry message, request_id |
| Seek hangs | Wait for buffering, check Range/proxy | Upstream status and proxy mode |
| Multiple users fail | Ask an admin to inspect WebSocket, Ingress, or Provider | Time range, room ID, affected users |

@ -16,7 +16,7 @@ import { Aside, Steps } from '@astrojs/starlight/components';
2. SyncTV 根据房间状态、媒体、用户、Provider 凭据和客户端能力生成 `Playback`。
3. Provider 返回一个或多个 `PlaybackInfo`,每个模式包含 URL、format、header、字幕、过期时间和 metadata。
4. 客户端选择合适模式:直连、proxy、转码、HLS、FLV 或字幕 URL。
5. URL 过期、媒体切换、Provider 凭据变化或客户端能力变化后,客户端重新获取快照。
5. URL 过期、媒体切换、Provider 凭据变化或客户端能力变化后,客户端重新获取播放信息。
</Steps>
## 直连和代理
@ -41,7 +41,7 @@ SyncTV proxy 只使用 Provider 明确给出的 header。它不会把浏览器
| 浏览器不能设置 `Referer` 或自定义 header | 优先选择 proxy URL |
| 原生客户端可设置 header 且能访问上游 | 可使用直连 URL |
| 移动端不支持某些 codec/container | 在 `PlaybackClientProfile` 中声明能力,选择转码或兼容模式 |
| URL 有 `expires_at` | 到期前刷新快照,不要无限重试旧 URL |
| URL 有 `expires_at` | 到期前重新获取播放信息 |
| 字幕有独立 header | 使用字幕 URL 自带 header;proxy 字幕时 header 会合并 |
## Range 与 slice cache
@ -69,7 +69,7 @@ SyncTV proxy slice cache 的边界:
- `playback_state`:当前播放位置、状态和版本。
- `playback`:当前媒体的播放 URL、headers、字幕和过期时间。
断线重连时,带上本地 `version`、`media_id`、`playlist_id`、`target` 和 `playback_client_profile`。服务端会判断缓存是否仍有效,必要时推送新的快照。
断线重连时,重新获取 `playback_state`,并用当前 `playback_client_profile` 观察 `playback`。播放信息是当前播放源可直接交给播放器使用的数据,生命周期跟随其中的 URL。
完整协议见 [Realtime API](../../develop/realtime-api/)。

@ -867,7 +867,7 @@ fn register_read_routes(_state: &AppState) -> Router<AppState> {
get(room::watch_playback_state),
)
.route(
"/api/rooms/{room_id}/watch/playback-snapshot",
"/api/rooms/{room_id}/watch/playback",
get(room::watch_playback),
)
.route(

@ -3817,8 +3817,8 @@ mod tests {
sse_event_from_server_message, sse_event_id_from_resource_changed,
watch_after_event_sequence, AddMediaBatchBody, CancelOnDropStream, ChatImageObjectQuery,
CreatePlaylistBody, DeleteEntriesBody, GetPlaybackQuery, PlaylistCoverObjectQuery,
RoomCoverObjectQuery, UpdatePlaybackRequest, VideoCoverObjectQuery,
WatchPlaybackQuery, WatchQuery,
RoomCoverObjectQuery, UpdatePlaybackRequest, VideoCoverObjectQuery, WatchPlaybackQuery,
WatchQuery,
};
use crate::proto::client::{
DeleteMediaQuery, DeletePlaylistQuery, GetChatHistoryRequest, GetChatMessageContextRequest,
@ -4007,6 +4007,10 @@ mod tests {
"format=json&media_id=media_1&extra=true"
)
.is_err());
assert!(serde_urlencoded::from_str::<WatchPlaybackQuery>(
"format=json&after_event_sequence=12"
)
.is_err());
assert!(
serde_urlencoded::from_str::<ChatImageObjectQuery>("token=token&extra=true").is_err()
);

@ -34,7 +34,7 @@ use super::client::convert::{
json_to_vec, playback_client_profile_from_proto, provider_playback_info_to_model,
sign_local_bilibili_danmaku_urls, try_bilibili_live_danmaku_for_static_media,
try_media_list_to_proto, try_media_to_proto, try_media_to_proto_with_availability,
try_members_to_proto, try_playback_to_proto, try_playback_state_to_proto,
try_members_to_proto, try_playback_state_to_proto, try_playback_to_proto,
try_playlist_list_to_proto, try_playlist_path_node_to_proto, try_playlist_to_proto,
try_playlist_to_proto_with_availability, user_status_to_proto,
};
@ -9992,9 +9992,7 @@ mod tests {
public_media_id(&admin_api, media.id)
);
let result = response
.playback
.expect("playback should be present");
let result = response.playback.expect("playback should be present");
assert_eq!(result.media_id, public_media_id(&admin_api, media.id));
assert_eq!(result.room_id, public_room_id(&admin_api, room.id));
assert_eq!(result.name, media.name);
@ -10002,7 +10000,8 @@ mod tests {
#[tokio::test]
#[ignore = "Requires Docker"]
async fn test_get_playback_returns_state_when_snapshot_generation_fails_for_global_admin() {
async fn test_get_playback_returns_state_when_playback_info_generation_fails_for_global_admin()
{
let (_postgres, pool) = create_test_pool().await;
let (admin_api, _redis_publish_rx) =
make_admin_api_for_delete_user_test(pool.clone()).await;
@ -10054,7 +10053,9 @@ mod tests {
let response = admin_api
.get_playback(&public_room_id(&admin_api, room.id), &global_admin.id, None)
.await
.expect("global admin should get playback state even if snapshot generation fails");
.expect(
"global admin should get playback state even if playback info generation fails",
);
let state = response
.playback_state
@ -10066,7 +10067,7 @@ mod tests {
);
assert!(
response.playback.is_none(),
"admin playback queries should degrade to state-only responses on snapshot failures"
"admin playback queries should degrade to state-only responses on playback info failures"
);
}
@ -10134,9 +10135,7 @@ mod tests {
.await
.expect("global admin should get signed provider playback");
let result = response
.playback
.expect("playback should be present");
let result = response.playback.expect("playback should be present");
let direct = result
.playback_infos
.get("direct")
@ -10227,9 +10226,7 @@ mod tests {
.await
.expect("local management actor should get signed provider playback");
let result = response
.playback
.expect("playback should be present");
let result = response.playback.expect("playback should be present");
let direct = result
.playback_infos
.get("direct")

@ -338,6 +338,12 @@ impl ResourceObservation {
}
}
fn exposes_client_event_cursor(&self) -> bool {
!matches!(self.resource, ObservedResource::Playback { .. })
&& (matches!(self.resource, ObservedResource::ChatEvents)
|| self.room_resource_cursor_types().is_some())
}
fn evaluation_key(&self) -> ObservationEvaluationKey {
let delivery_mode = self.delivery_mode as i32;
match &self.resource {
@ -615,9 +621,7 @@ impl ResourceObserver {
request: &crate::proto::client::ObserveResource,
) -> i64 {
let requested_sequence = Self::requested_replay_sequence(request).unwrap_or(0).max(0);
if matches!(observation.resource, ObservedResource::ChatEvents)
|| observation.room_resource_cursor_types().is_some()
{
if observation.exposes_client_event_cursor() {
requested_sequence
} else {
0
@ -1292,7 +1296,9 @@ impl MediaResourceHub {
if let (Some(changed), Some(cursor)) =
(changed_message.as_mut(), event_cursor.as_ref())
{
changed.event_cursor = Some(cursor.clone());
if entry.observation.exposes_client_event_cursor() {
changed.event_cursor = Some(cursor.clone());
}
}
match self
.send_and_commit_subscription_update(
@ -1468,7 +1474,6 @@ impl ResourceObserver {
}
fn observation_from_request(
&self,
request: &crate::proto::client::ObserveResource,
) -> Result<ResourceObservation, String> {
use crate::proto::client::observe_resource::Resource;
@ -1498,17 +1503,17 @@ impl ResourceObserver {
}
Resource::RoomSettings(_) => ObservedResource::RoomSettings,
Resource::PlaylistItems(observe) => ObservedResource::PlaylistItems {
request: observe
.request
.clone()
.ok_or_else(|| "playlist_items request is required".to_string())?,
},
request: observe
.request
.clone()
.ok_or_else(|| "playlist_items request is required".to_string())?,
},
Resource::RoomMembers(observe) => ObservedResource::RoomMembers {
request: observe
.request
.clone()
.ok_or_else(|| "room_members request is required".to_string())?,
},
request: observe
.request
.clone()
.ok_or_else(|| "room_members request is required".to_string())?,
},
Resource::ChatEvents(_) => ObservedResource::ChatEvents,
};
@ -1526,7 +1531,7 @@ impl ResourceObserver {
self: &Arc<Self>,
request: &crate::proto::client::ObserveResource,
) -> Result<(), String> {
let mut observation = match self.observation_from_request(request) {
let mut observation = match Self::observation_from_request(request) {
Ok(observation) => observation,
Err(error) => {
self.send_server_message(Self::resource_observe_error_message(
@ -1552,7 +1557,8 @@ impl ResourceObserver {
let start_sequence = Self::observation_start_sequence(&observation, request);
let is_chat_observation = matches!(observation.resource, ObservedResource::ChatEvents);
let observed_cursor = if is_chat_observation {
let exposes_client_event_cursor = observation.exposes_client_event_cursor();
let internal_cursor = if is_chat_observation {
Some(crate::proto::client::EventCursor {
event_id: None,
sequence: start_sequence,
@ -1569,9 +1575,12 @@ impl ResourceObserver {
_ => None,
}
};
observation.last_sent_event_sequence = observed_cursor
observation.last_sent_event_sequence = internal_cursor
.as_ref()
.map_or(start_sequence, |cursor| cursor.sequence);
let observed_cursor = internal_cursor
.clone()
.filter(|_| exposes_client_event_cursor);
match self.evaluate_observation(&mut observation).await {
Ok(mut update) => {
if is_chat_observation {
@ -1729,6 +1738,9 @@ impl ResourceObserver {
let Some(observation) = self.local_observation(observe_id).await else {
return Ok(());
};
if !observation.exposes_client_event_cursor() {
return Ok(());
}
let Some(resource_types) = observation.room_resource_cursor_types() else {
return Ok(());
};
@ -1852,9 +1864,7 @@ impl ResourceObserver {
Ok(())
}
pub(super) async fn next_playback_refresh_deadline(
&self,
) -> Option<tokio::time::Instant> {
pub(super) async fn next_playback_refresh_deadline(&self) -> Option<tokio::time::Instant> {
let state = self.state.lock().await;
let expires_at = state
.observations
@ -1877,9 +1887,7 @@ impl ResourceObserver {
)
}
pub(super) async fn refresh_expired_playback_observations(
&self,
) -> Result<(), String> {
pub(super) async fn refresh_expired_playback_observations(&self) -> Result<(), String> {
self.room_hub
.refresh_expired_playbacks(Some(&self.connection_id))
.await
@ -2023,10 +2031,7 @@ impl ResourceObserver {
_observation: &ResourceObservation,
invalidation: &ResourceInvalidation,
) -> bool {
match invalidation {
ResourceInvalidation::Playback(_) => true,
_ => false,
}
matches!(invalidation, ResourceInvalidation::Playback(_))
}
async fn evaluate_observation(
@ -2132,18 +2137,16 @@ impl ResourceObserver {
)),
}
}
ObservedResource::Playback { .. } => {
self.playback_service.as_ref().map_or(
SharedResourceServiceIdentity {
id: 0,
weak: SharedResourceServiceWeak::AlwaysAlive,
},
|service| SharedResourceServiceIdentity {
id: Arc::as_ptr(service).cast::<()>() as usize,
weak: SharedResourceServiceWeak::Playback(Arc::downgrade(service)),
},
)
}
ObservedResource::Playback { .. } => self.playback_service.as_ref().map_or(
SharedResourceServiceIdentity {
id: 0,
weak: SharedResourceServiceWeak::AlwaysAlive,
},
|service| SharedResourceServiceIdentity {
id: Arc::as_ptr(service).cast::<()>() as usize,
weak: SharedResourceServiceWeak::Playback(Arc::downgrade(service)),
},
),
ObservedResource::RoomSettings => {
let id = self.room_settings_snapshot_service_id;
SharedResourceServiceIdentity {
@ -2238,11 +2241,7 @@ impl ResourceObserver {
Payload::Playback(playback.clone())
}
};
(
hex::encode(fingerprint),
playback.expires_at,
payload,
)
(hex::encode(fingerprint), playback.expires_at, payload)
}
ObservedResource::RoomSettings => {
let service = Arc::clone(&self.room_settings_snapshot_service);
@ -2365,7 +2364,7 @@ mod tests {
fn playback_observation() -> ResourceObservation {
ResourceObservation {
observe_id: "playback-snapshot".to_string(),
observe_id: "playback".to_string(),
last_fingerprint: String::new(),
delivery_mode: ResourceDeliveryMode::PushSnapshot,
resource: ObservedResource::Playback {
@ -2377,7 +2376,7 @@ mod tests {
}
#[test]
fn playback_uses_room_resource_event_cursor_types() {
fn playback_tracks_room_resource_dependency_types() {
let observation = playback_observation();
assert_eq!(
@ -2387,31 +2386,26 @@ mod tests {
}
#[test]
fn playback_is_room_resource_cursor_observation() {
fn playback_hides_client_event_cursor() {
let observation = playback_observation();
assert!(observation.room_resource_cursor_types().is_some());
assert!(!observation.exposes_client_event_cursor());
}
#[test]
fn playback_observation_ignores_replay_sequence() {
fn playback_observation_has_no_requested_event_sequence() {
let observation = playback_observation();
let request = crate::proto::client::ObserveResource {
observe_id: "playback".to_string(),
delivery_mode: ResourceDeliveryMode::PushSnapshot as i32,
resource: Some(
crate::proto::client::observe_resource::Resource::Playback(
crate::proto::client::ObservePlayback {
playback_client_profile: None,
},
),
),
resource: Some(crate::proto::client::observe_resource::Resource::Playback(
crate::proto::client::ObservePlayback {
playback_client_profile: None,
},
)),
};
assert_eq!(
ResourceObserver::requested_replay_sequence(&request),
None
);
assert_eq!(ResourceObserver::requested_replay_sequence(&request), None);
assert_eq!(
ResourceObserver::observation_start_sequence(&observation, &request),
0
@ -2419,18 +2413,16 @@ mod tests {
}
#[test]
fn playback_observation_uses_latest_room_resource_cursor_for_live_start() {
fn playback_observation_starts_from_current_playback_without_client_event_cursor() {
let observation = playback_observation();
let request = crate::proto::client::ObserveResource {
observe_id: "playback".to_string(),
delivery_mode: ResourceDeliveryMode::PushSnapshot as i32,
resource: Some(
crate::proto::client::observe_resource::Resource::Playback(
crate::proto::client::ObservePlayback {
playback_client_profile: None,
},
),
),
resource: Some(crate::proto::client::observe_resource::Resource::Playback(
crate::proto::client::ObservePlayback {
playback_client_profile: None,
},
)),
};
assert_eq!(
@ -2440,7 +2432,7 @@ mod tests {
assert_eq!(
ResourceObserver::requested_replay_sequence(&request),
None,
"absent playback cursor starts live and uses the latest durable room cursor"
"playback observation has no client event cursor"
);
}

File diff suppressed because it is too large Load Diff

@ -42,10 +42,8 @@ pub trait PlaybackService: Send + Sync {
}
}
pub(crate) fn playback_expires_at(
snapshot: &crate::proto::client::Playback,
) -> Option<i64> {
snapshot
pub(crate) fn playback_expires_at(playback: &crate::proto::client::Playback) -> Option<i64> {
playback
.playback_infos
.values()
.flat_map(|info| info.urls.iter().filter_map(|url| url.expire_at))

@ -65,9 +65,7 @@ pub fn resource_invalidations_for_room_event(event: &RealtimeEvent) -> Vec<Resou
match event {
RealtimeEvent::PlaybackStateChanged { .. } => vec![
ResourceInvalidation::PlaybackState,
ResourceInvalidation::Playback(
PlaybackInvalidation::PlaybackStateChanged,
),
ResourceInvalidation::Playback(PlaybackInvalidation::PlaybackStateChanged),
],
RealtimeEvent::MediaUpdated { media_id, .. }
| RealtimeEvent::MediaRemoved { media_id, .. } => vec![
@ -85,11 +83,9 @@ pub fn resource_invalidations_for_room_event(event: &RealtimeEvent) -> Vec<Resou
RealtimeEvent::MediaRemovedBatch { media_ids, .. }
| RealtimeEvent::PlaylistReordered { media_ids, .. } => vec![
ResourceInvalidation::PlaylistItems,
ResourceInvalidation::Playback(
PlaybackInvalidation::PlaylistItemsChanged {
media_ids: media_ids.clone(),
},
),
ResourceInvalidation::Playback(PlaybackInvalidation::PlaylistItemsChanged {
media_ids: media_ids.clone(),
}),
],
RealtimeEvent::PlaylistDeleted { playlist_id, .. } => vec![
ResourceInvalidation::PlaylistItems,
@ -235,7 +231,7 @@ mod tests {
}
#[test]
fn playback_state_event_invalidates_state_and_snapshot() {
fn playback_state_event_invalidates_state_and_playback() {
let event = RealtimeEvent::PlaybackStateChanged {
event_id: "evt".to_string(),
room_id: room_id(),
@ -249,9 +245,7 @@ mod tests {
resource_invalidations_for_room_event(&event),
vec![
ResourceInvalidation::PlaybackState,
ResourceInvalidation::Playback(
PlaybackInvalidation::PlaybackStateChanged
),
ResourceInvalidation::Playback(PlaybackInvalidation::PlaybackStateChanged),
]
);
}
@ -310,9 +304,9 @@ mod tests {
resource_invalidations_for_room_event(&event),
vec![
ResourceInvalidation::PlaylistItems,
ResourceInvalidation::Playback(
PlaybackInvalidation::PlaylistChanged { playlist_id }
),
ResourceInvalidation::Playback(PlaybackInvalidation::PlaylistChanged {
playlist_id
}),
]
);
}
@ -333,9 +327,7 @@ mod tests {
resource_invalidations_for_room_event(&event),
vec![
ResourceInvalidation::PlaylistItems,
ResourceInvalidation::Playback(
PlaybackInvalidation::MediaChanged { media_id }
),
ResourceInvalidation::Playback(PlaybackInvalidation::MediaChanged { media_id }),
]
);
}
@ -361,9 +353,9 @@ mod tests {
};
let expected = vec![
ResourceInvalidation::PlaylistItems,
ResourceInvalidation::Playback(
PlaybackInvalidation::PlaylistItemsChanged { media_ids },
),
ResourceInvalidation::Playback(PlaybackInvalidation::PlaylistItemsChanged {
media_ids,
}),
];
assert_eq!(resource_invalidations_for_room_event(&removed), expected);
@ -386,9 +378,9 @@ mod tests {
resource_invalidations_for_room_event(&event),
vec![
ResourceInvalidation::PlaylistItems,
ResourceInvalidation::Playback(
PlaybackInvalidation::PlaylistChanged { playlist_id }
),
ResourceInvalidation::Playback(PlaybackInvalidation::PlaylistChanged {
playlist_id
}),
]
);
}

@ -541,12 +541,12 @@ async fn test_get_playback_returns_dynamic_playlist_item_playback_info() {
update_state.position < 17.5,
"playing update response should not jump far beyond the requested seek position"
);
let update_snapshot = update_response
let update_playback = update_response
.playback
.expect("update response should preserve read-after-write playback");
assert_eq!(update_snapshot.playlist_id, playlist_public_id);
assert_eq!(update_snapshot.name, "episode-1.mp4");
let update_direct = update_snapshot.playback_infos.get("direct").unwrap();
assert_eq!(update_playback.playlist_id, playlist_public_id);
assert_eq!(update_playback.name, "episode-1.mp4");
let update_direct = update_playback.playback_infos.get("direct").unwrap();
assert_eq!(
update_direct.urls[0].url,
"https://alist_default.example.com/episode-1.mp4"
@ -892,7 +892,7 @@ async fn test_static_provider_playback_with_signing_key_uses_provider_store_regi
#[tokio::test]
#[ignore = "Requires Docker"]
async fn test_get_playback_without_active_media_returns_stable_non_empty_snapshot_version() {
async fn test_get_playback_without_active_media_returns_idle_playback_info() {
let (_postgres, pool) = create_test_pool().await;
let user_repo = UserRepository::new(pool.clone());
let user_service = Arc::new(make_user_service(&pool));
@ -963,7 +963,7 @@ async fn test_get_playback_without_active_media_returns_stable_non_empty_snapsho
#[tokio::test]
#[ignore = "Requires Docker"]
async fn test_get_playback_returns_state_when_snapshot_generation_fails() {
async fn test_get_playback_returns_state_when_playback_info_generation_fails() {
let (_postgres, pool) = create_test_pool().await;
let user_repo = UserRepository::new(pool.clone());
let media_repo = MediaRepository::new(pool.clone());
@ -1051,12 +1051,12 @@ async fn test_get_playback_returns_state_when_snapshot_generation_fails() {
let playback_state = response
.playback_state
.expect("playback state should still be returned when snapshot generation fails");
.expect("playback state should still be returned when playback info generation fails");
assert_eq!(playback_state.playing_media_id, media_public_id);
assert!(playback_state.is_playing);
assert!(
response.playback.is_none(),
"snapshot failures should degrade to state-only responses"
"playback info failures should degrade to state-only responses"
);
}

@ -6637,13 +6637,13 @@ fn build_get_playback_cli_output(
let mut flv_pull_url = None;
let mut flv_absolute_pull_url = None;
if let Some(snapshot) = playback.as_ref() {
let mut modes = snapshot.playback_infos.iter().collect::<Vec<_>>();
if let Some(playback) = playback.as_ref() {
let mut modes = playback.playback_infos.iter().collect::<Vec<_>>();
modes.sort_by_key(|(mode, _)| *mode);
for (mode, info) in modes {
for (index, playback_url) in info.urls.iter().enumerate() {
let is_default = mode == &snapshot.default_mode
let is_default = mode == &playback.default_mode
&& i32::try_from(index).is_ok_and(|index| info.default_url_index == index);
let absolute_url = absolutize_cli_url(&playback_url.url, api_base_url.as_deref());
let output = PlaybackPullUrlCliOutput {
@ -6681,7 +6681,7 @@ fn build_get_playback_cli_output(
let default_mode = playback
.as_ref()
.map(|snapshot| snapshot.default_mode.clone())
.map(|playback| playback.default_mode.clone())
.filter(|mode| !mode.is_empty());
GetPlaybackCliOutput {
@ -13664,7 +13664,6 @@ mod tests {
)]),
default_mode: "direct".into(),
metadata: std::collections::HashMap::new(),
version: "1".into(),
expires_at: None,
}),
},

@ -115,28 +115,18 @@ fn optional_event_sequence(sequence: impl Into<String>) -> Option<i64> {
}
}
fn observe_playback_message(
observe_id: &str,
after_event_sequence: impl Into<String>,
) -> synctv_proto::client::ClientMessage {
fn observe_playback_message(observe_id: &str) -> synctv_proto::client::ClientMessage {
use synctv_proto::client::{
client_message, observe_resource, ObservePlayback, ObserveResource,
ResourceDeliveryMode,
client_message, observe_resource, ObservePlayback, ObserveResource, ResourceDeliveryMode,
};
synctv_proto::client::ClientMessage {
message: Some(client_message::Message::ObserveResource(ObserveResource {
observe_id: observe_id.to_string(),
delivery_mode: ResourceDeliveryMode::PushSnapshot as i32,
resource: Some(observe_resource::Resource::Playback(
ObservePlayback {
media_id: None,
playlist_id: None,
target: Vec::new(),
playback_client_profile: None,
after_event_sequence: optional_event_sequence(after_event_sequence),
},
)),
resource: Some(observe_resource::Resource::Playback(ObservePlayback {
playback_client_profile: None,
})),
})),
}
}
@ -210,9 +200,7 @@ fn observe_room_members_message(
}
}
fn resource_playback(
message: &ServerMessage,
) -> Option<&synctv_proto::client::Playback> {
fn resource_playback(message: &ServerMessage) -> Option<&synctv_proto::client::Playback> {
match &message.message {
Some(server_message::Message::ResourceChanged(changed)) => match changed.payload.as_ref() {
Some(synctv_proto::client::resource_changed::Payload::Playback(snapshot)) => {
@ -4215,10 +4203,7 @@ async fn full_stack_cli_media_resource_and_member_commands_cover_status_permissi
media_one_id
);
assert_eq!(playback_state["playback_state"]["is_playing"], true);
assert_eq!(
playback_state["playback"]["media_id"],
media_one_id
);
assert_eq!(playback_state["playback"]["media_id"], media_one_id);
let stopped_playback = run_synctv_remote_cli_json(
&server,
@ -7073,12 +7058,11 @@ async fn full_stack_grpc_message_stream_establishes_and_acks_heartbeat() {
#[tokio::test]
#[ignore = "Requires Docker (testcontainers)"]
async fn full_stack_grpc_message_stream_watch_playback_receives_initial_and_future_updates(
) {
async fn full_stack_grpc_message_stream_watch_playback_receives_initial_and_future_updates() {
use synctv_proto::client::room_service_client::RoomServiceClient;
use tokio_stream::wrappers::ReceiverStream;
let fixture = start_room_realtime_fixture("grpc-watch-playback-snapshot").await;
let fixture = start_room_realtime_fixture("grpc-watch-playback").await;
let RoomRealtimeFixture {
server,
room_id,
@ -7149,10 +7133,7 @@ async fn full_stack_grpc_message_stream_watch_playback_receives_initial_and_futu
.expect("connect member room gRPC client");
let (outbound_tx, outbound_rx) = tokio::sync::mpsc::channel(8);
outbound_tx
.send(observe_playback_message(
"grpc-playback-snapshot",
String::new(),
))
.send(observe_playback_message("grpc-playback"))
.await
.expect("queue initial playback observe request");
let outbound = ReceiverStream::new(outbound_rx);
@ -7169,21 +7150,15 @@ async fn full_stack_grpc_message_stream_watch_playback_receives_initial_and_futu
.expect("message_stream should establish")
.into_inner();
let initial_snapshot = recv_matching_grpc_server_message(
let _initial_playback = recv_matching_grpc_server_message(
&mut inbound,
Duration::from_secs(10),
|message| {
resource_playback(message).is_some_and(|snapshot| {
snapshot.media_id == media_one_id && !snapshot.version.is_empty()
})
resource_playback(message).is_some_and(|playback| playback.media_id == media_one_id)
},
"initial grpc playback",
)
.await;
let initial_version = resource_playback(&initial_snapshot)
.expect("playback should be present")
.version
.clone();
let _ = run_synctv_remote_cli_json(
&server,
@ -7200,25 +7175,19 @@ async fn full_stack_grpc_message_stream_watch_playback_receives_initial_and_futu
)
.await;
let updated_snapshot = recv_matching_grpc_server_message(
let updated_playback = recv_matching_grpc_server_message(
&mut inbound,
Duration::from_secs(10),
|message| {
resource_playback(message).is_some_and(|snapshot| {
snapshot.media_id == media_two_id && !snapshot.version.is_empty()
})
resource_playback(message).is_some_and(|playback| playback.media_id == media_two_id)
},
"updated grpc playback",
)
.await;
let snapshot =
resource_playback(&updated_snapshot).expect("updated snapshot should be present");
assert_eq!(snapshot.media_id, media_two_id);
assert_eq!(
snapshot.version, initial_version,
"switching between newly-created media should preserve resource version"
);
let playback =
resource_playback(&updated_playback).expect("updated playback should be present");
assert_eq!(playback.media_id, media_two_id);
}
#[tokio::test]
@ -8341,7 +8310,7 @@ async fn full_stack_websocket_room_messages_include_playlist_lifecycle_events()
async fn full_stack_websocket_watch_playback_receives_initial_and_future_updates() {
use synctv_proto::client::server_message;
let fixture = start_room_realtime_fixture("ws-watch-playback-snapshot").await;
let fixture = start_room_realtime_fixture("ws-watch-playback").await;
let RoomRealtimeFixture {
server,
api_addr,
@ -8464,27 +8433,17 @@ async fn full_stack_websocket_watch_playback_receives_initial_and_future_updates
)
.await;
send_client_message(
&mut member_ws,
observe_playback_message("ws-playback-snapshot", String::new()),
)
.await;
send_client_message(&mut member_ws, observe_playback_message("ws-playback")).await;
let initial_snapshot = recv_matching_server_message(
let _initial_playback = recv_matching_server_message(
&mut member_ws,
Duration::from_secs(10),
|message| {
resource_playback(message).is_some_and(|snapshot| {
snapshot.media_id == media_one_id && !snapshot.version.is_empty()
})
resource_playback(message).is_some_and(|playback| playback.media_id == media_one_id)
},
"initial playback",
)
.await;
let initial_version = resource_playback(&initial_snapshot)
.expect("playback should be present")
.version
.clone();
let _ = run_synctv_remote_cli_json(
&server,
@ -8501,25 +8460,19 @@ async fn full_stack_websocket_watch_playback_receives_initial_and_future_updates
)
.await;
let updated_snapshot = recv_matching_server_message(
let updated_playback = recv_matching_server_message(
&mut member_ws,
Duration::from_secs(10),
|message| {
resource_playback(message).is_some_and(|snapshot| {
snapshot.media_id == media_two_id && !snapshot.version.is_empty()
})
resource_playback(message).is_some_and(|playback| playback.media_id == media_two_id)
},
"updated playback",
)
.await;
let snapshot =
resource_playback(&updated_snapshot).expect("updated snapshot should be present");
assert_eq!(snapshot.media_id, media_two_id);
assert_eq!(
snapshot.version, initial_version,
"switching between newly-created media should preserve resource version"
);
let playback =
resource_playback(&updated_playback).expect("updated playback should be present");
assert_eq!(playback.media_id, media_two_id);
}
#[tokio::test]

Loading…
Cancel
Save