Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions Caddyfile
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,17 @@ tracker.rocketup.tech {
}
}

@sse path /api/events /api/events/*
handle @sse {
reverse_proxy backend:3002 {
flush_interval -1
transport http {
read_timeout 0
write_timeout 0
}
}
}

handle /api/* {
reverse_proxy backend:3002
}
Expand Down
112 changes: 112 additions & 0 deletions alfy-bot-frontend/src/composables/useSse.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,112 @@
import { onUnmounted } from 'vue'
import { api } from '@/api/client'
import { getToken } from '@/api/tokenStorage'

export type SseHandler = (event: string, data: unknown) => void

const INITIAL_RETRY_MS = 1000
const MAX_RETRY_MS = 30_000

function parseSseBlock(block: string): { event: string, data: string } | null {
let event = 'message'
let data = ''
for (const line of block.split('\n')) {
if (line.startsWith('event:'))
event = line.slice(6).trim()
else if (line.startsWith('data:'))
data += line.slice(5).trim()
}
if (!data)
return null
return { event, data }
}

export function useSse(path: string, onEvent: SseHandler) {
let abort: AbortController | null = null
let retryMs = INITIAL_RETRY_MS
let stopped = false
let retryTimer: ReturnType<typeof setTimeout> | null = null

async function readStream(body: ReadableStream<Uint8Array>) {
const reader = body.getReader()
const decoder = new TextDecoder()
let buf = ''

while (true) {
if (stopped)
break
const { done, value } = await reader.read()
if (done)
break
buf += decoder.decode(value, { stream: true })
const parts = buf.split('\n\n')
buf = parts.pop() ?? ''
for (const block of parts) {
const parsed = parseSseBlock(block)
if (!parsed)
continue
try {
onEvent(parsed.event, JSON.parse(parsed.data))
}
catch {
onEvent(parsed.event, parsed.data)
}
}
}
}

async function connect() {
const token = getToken()
if (!token)
return

abort = new AbortController()
const base = String(api.defaults.baseURL ?? '').replace(/\/$/, '')
const res = await fetch(`${base}${path}`, {
headers: {
Authorization: `Bearer ${token}`,
Accept: 'text/event-stream',
},
signal: abort.signal,
})

if (res.status === 401) {
stopped = true
return
}
if (!res.ok || !res.body)
throw new Error(`SSE ${res.status}`)

retryMs = INITIAL_RETRY_MS
await readStream(res.body)
}

async function loop() {
while (true) {
if (stopped)
return
try {
await connect()
}
catch (e) {
if (stopped || (e instanceof DOMException && e.name === 'AbortError'))
return
}
if (stopped)
return
await new Promise<void>((resolve) => {
retryTimer = setTimeout(resolve, retryMs)
})
retryMs = Math.min(retryMs * 2, MAX_RETRY_MS)
}
}

loop()

onUnmounted(() => {
stopped = true
if (retryTimer)
clearTimeout(retryTimer)
abort?.abort()
})
}
32 changes: 26 additions & 6 deletions alfy-bot-frontend/src/features/task-timer/model/timer-store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,11 @@ export const useTimerStore = defineStore('timer', () => {
const timerInterval: Ref<ReturnType<typeof setInterval> | null> = ref(null)
const { play: playSound } = useSounds()
let swListenerRegistered = false
let sseMutedUntil = 0

function muteRemoteRestore(ms = 1500): void {
sseMutedUntil = Date.now() + ms
}

const checkPhase = {
isWorkPhase: (phaseNumber: number): boolean => phaseNumber % 2 === 1,
Expand Down Expand Up @@ -200,6 +205,8 @@ export const useTimerStore = defineStore('timer', () => {
const taskId = currentSettings.value.taskId
if (!taskId) return

muteRemoteRestore()

try {
if (isActive.value) {
const expiresAt = new Date(Date.now() + timeBlock.value * 1000).toISOString()
Expand Down Expand Up @@ -229,9 +236,14 @@ export const useTimerStore = defineStore('timer', () => {
async function restoreSession(): Promise<void> {
registerSWListener()

if (Date.now() < sseMutedUntil) return

try {
const { data: session } = await api.get('/tasks/timer')
if (!session) return
if (!session) {
clearLocalSession()
return
}

if (session.task?.pomodoroConfig) {
const cfg = session.task.pomodoroConfig
Expand All @@ -245,6 +257,8 @@ export const useTimerStore = defineStore('timer', () => {
}
}

stopLocalTicker()

const sessionState = determineSessionState(session)

switch (sessionState.type) {
Expand Down Expand Up @@ -318,28 +332,34 @@ export const useTimerStore = defineStore('timer', () => {
}
}

async function resetToInitialState(): Promise<void> {
function stopLocalTicker(): void {
if (timerInterval.value) {
clearInterval(timerInterval.value)
timerInterval.value = null
}

const timerId = currentSettings.value.taskId || 'pomodoro'
sendToSW({ type: 'TIMER_STOP', data: { id: timerId } })

isActive.value = false
expiresAt.value = null
}

function clearLocalSession(): void {
stopLocalTicker()
phase.value = 0
timeBlock.value = 0
namePhase.value = ''
updateTitle()
}

async function resetToInitialState(): Promise<void> {
muteRemoteRestore()
clearLocalSession()

try {
await api.delete('/tasks/timer')
} catch (err) {
console.error('Ошибка деактивации таймера:', err)
}

updateTitle()
}

function startTask(task: Task): void {
Expand Down
10 changes: 8 additions & 2 deletions alfy-bot-frontend/src/features/task-timer/ui/TimeBlock.vue
Original file line number Diff line number Diff line change
@@ -1,8 +1,9 @@
<script setup lang="ts">
import { computed, onMounted, onUnmounted, ref, watch } from 'vue'
import { useTimerStore } from '../model/timer-store'
import { Pause, Play, Square } from 'lucide-vue-next'
import { Button } from '@/components/ui/button'
import { Play, Pause, Square } from 'lucide-vue-next'
import { useSse } from '@/composables/useSse'
import { useTimerStore } from '../model/timer-store'

const store = useTimerStore()
const isFullscreen = ref(false)
Expand Down Expand Up @@ -56,6 +57,11 @@ onMounted(() => {
store.restoreSession()
})

useSse('/events', (event) => {
if (event === 'timer.updated')
void store.restoreSession()
})

onUnmounted(() => {
releaseWakeLock()
})
Expand Down
4 changes: 3 additions & 1 deletion alfy-bot-frontend/src/sw.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,9 @@ cleanupOutdatedCaches()

registerRoute(
({ url }) =>
url.origin === self.location.origin && url.pathname.startsWith('/api'),
url.origin === self.location.origin &&
url.pathname.startsWith('/api') &&
!url.pathname.startsWith('/api/events'),
new NetworkOnly(),
)

Expand Down
108 changes: 108 additions & 0 deletions alfy-bot-frontend/tests/composables/useSse.spec.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,108 @@
import type { VueWrapper } from '@vue/test-utils'
import { mount } from '@vue/test-utils'
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
import { defineComponent } from 'vue'
import { getToken } from '@/api/tokenStorage'
import { useSse } from '@/composables/useSse'

vi.mock('@/api/tokenStorage', () => ({
getToken: vi.fn(),
}))

vi.mock('@/api/client', () => ({
api: { defaults: { baseURL: 'http://localhost:3002/api' } },
}))

function sseStream(chunks: string[], hang = true): ReadableStream<Uint8Array> {
const encoder = new TextEncoder()
return new ReadableStream({
start(controller) {
for (const chunk of chunks) {
controller.enqueue(encoder.encode(chunk))
}
if (!hang)
controller.close()
},
})
}

describe('useSse', () => {
let wrapper: VueWrapper | null = null
let fetchMock: ReturnType<typeof vi.fn>

beforeEach(() => {
vi.mocked(getToken).mockReturnValue('jwt-token')
fetchMock = vi.fn()
vi.stubGlobal('fetch', fetchMock)
})

afterEach(() => {
wrapper?.unmount()
wrapper = null
vi.unstubAllGlobals()
})

function mountSse(onEvent: (event: string, data: unknown) => void) {
wrapper = mount(defineComponent({
setup() {
useSse('/events', onEvent)
return () => null
},
}))
return wrapper
}

it('не коннектится без токена', async () => {
vi.mocked(getToken).mockReturnValue(null)
mountSse(vi.fn())
await Promise.resolve()
expect(fetchMock).not.toHaveBeenCalled()
})

it('шлёт Bearer и парсит timer.updated', async () => {
const onEvent = vi.fn()
fetchMock.mockResolvedValue({
ok: true,
status: 200,
body: sseStream(['event: timer.updated\ndata: {}\n\n']),
})

mountSse(onEvent)
await vi.waitFor(() => expect(onEvent).toHaveBeenCalledWith('timer.updated', {}))

expect(fetchMock).toHaveBeenCalledWith(
'http://localhost:3002/api/events',
expect.objectContaining({
headers: expect.objectContaining({
Authorization: 'Bearer jwt-token',
Accept: 'text/event-stream',
}),
}),
)
})

it('не вызывает handler на ping', async () => {
const onEvent = vi.fn()
fetchMock.mockResolvedValue({
ok: true,
status: 200,
body: sseStream(['event: ping\ndata: {}\n\n', 'event: timer.updated\ndata: {}\n\n']),
})

mountSse((event, data) => {
if (event === 'ping')
return
onEvent(event, data)
})
await vi.waitFor(() => expect(onEvent).toHaveBeenCalledWith('timer.updated', {}))
expect(onEvent).toHaveBeenCalledTimes(1)
})

it('не ретраит 401', async () => {
fetchMock.mockResolvedValue({ ok: false, status: 401, body: null })
mountSse(vi.fn())
await Promise.resolve()
await Promise.resolve()
expect(fetchMock).toHaveBeenCalledTimes(1)
})
})
Loading
Loading