1
// org.freedesktop.Secret.Service
2

            
3
use std::{
4
    collections::HashMap,
5
    sync::{
6
        Arc, RwLock,
7
        atomic::{AtomicU32, Ordering},
8
    },
9
};
10

            
11
use oo7::{
12
    Key, Secret,
13
    dbus::{
14
        Algorithm, ServiceError,
15
        api::{DBusSecretInner, Properties},
16
    },
17
    file::{Keyring, LockedKeyring, UnlockedKeyring},
18
};
19
use tokio::sync::Mutex;
20
use tokio_stream::StreamExt;
21
use zbus::{
22
    fdo::{RequestNameFlags, RequestNameReply},
23
    names::UniqueName,
24
    object_server::SignalEmitter,
25
    proxy::Defaults,
26
    zvariant::{ObjectPath, Optional, OwnedObjectPath, OwnedValue, Value},
27
};
28

            
29
#[cfg(any(feature = "gnome_native_crypto", feature = "gnome_openssl_crypto"))]
30
pub use crate::gnome::internal::InternalInterface;
31
#[cfg(any(feature = "plasma_native_crypto", feature = "plasma_openssl_crypto"))]
32
use crate::plasma::prompter::in_plasma_environment;
33
use crate::{
34
    collection::Collection,
35
    error::{Error, custom_service_error},
36
    migration::{self, PendingMigration},
37
    prompt::{Prompt, PromptAction, PromptRole},
38
    session::{PeerInfo, Session, SessionType},
39
};
40

            
41
const DEFAULT_COLLECTION_ALIAS_PATH: ObjectPath<'static> =
42
    ObjectPath::from_static_str_unchecked("/org/freedesktop/secrets/aliases/default");
43

            
44
/// Prompter type
45
#[derive(Clone, Copy, PartialEq, Eq)]
46
pub enum PrompterType {
47
    #[allow(clippy::upper_case_acronyms)]
48
    GNOME,
49
    Plasma,
50
    Cli,
51
}
52

            
53
#[derive(Clone)]
54
pub struct Service {
55
    // Properties
56
    pub(crate) collections: Arc<Mutex<HashMap<OwnedObjectPath, Collection>>>,
57
    // Other attributes
58
    connection: Arc<RwLock<Option<zbus::Connection>>>,
59
    // sessions mapped to their corresponding object path on the bus
60
    sessions: Arc<Mutex<HashMap<OwnedObjectPath, Session>>>,
61
    session_index: Arc<AtomicU32>,
62
    // prompts mapped to their corresponding object path on the bus
63
    prompts: Arc<Mutex<HashMap<OwnedObjectPath, Prompt>>>,
64
    prompt_index: Arc<AtomicU32>,
65
    // pending collection creations: prompt_path -> (label, alias)
66
    pending_collections: Arc<Mutex<HashMap<OwnedObjectPath, (String, String)>>>,
67
    // pending keyring migrations: name -> migration
68
    pub(crate) pending_migrations: Arc<Mutex<HashMap<String, PendingMigration>>>,
69
    // Data directory for keyrings (e.g., ~/.local/share or test temp dir)
70
    data_dir: std::path::PathBuf,
71
    // PAM socket path (None for tests that don't need PAM listener)
72
    pub(crate) pam_socket: Option<std::path::PathBuf>,
73
    // Override for prompter type (mainly for tests)
74
    pub(crate) prompter_type_override: Arc<Mutex<Option<PrompterType>>>,
75
}
76

            
77
#[zbus::interface(name = "org.freedesktop.Secret.Service")]
78
impl Service {
79
    #[zbus(out_args("output", "result"))]
80
29
    pub async fn open_session(
81
        &self,
82
        algorithm: Algorithm,
83
        input: Value<'_>,
84
        #[zbus(header)] header: zbus::message::Header<'_>,
85
        #[zbus(object_server)] object_server: &zbus::ObjectServer,
86
    ) -> Result<(OwnedValue, OwnedObjectPath), ServiceError> {
87
63
        let (public_key, aes_key) = match algorithm {
88
29
            Algorithm::Plain => (None, None),
89
            Algorithm::Encrypted => {
90
36
                let client_public_key = Key::try_from(input).map_err(|err| {
91
                    custom_service_error(&format!(
92
                        "Input Value could not be converted into a Key {err}."
93
                    ))
94
                })?;
95
36
                let private_key = Key::generate_private_key().map_err(|err| {
96
                    custom_service_error(&format!("Failed to generate private key {err}."))
97
                })?;
98
                (
99
35
                    Some(Key::generate_public_key(&private_key).map_err(|err| {
100
                        custom_service_error(&format!("Failed to generate public key {err}."))
101
                    })?),
102
16
                    Some(
103
30
                        Key::generate_aes_key(&private_key, &client_public_key).map_err(|err| {
104
                            custom_service_error(&format!("Failed to generate aes key {err}."))
105
                        })?,
106
                    ),
107
                )
108
            }
109
        };
110

            
111
64
        let sender = if let Some(s) = header.sender() {
112
            s.to_owned()
113
        } else {
114
            #[cfg(any(test, feature = "test-util"))]
115
            {
116
                // For p2p test connections, use a dummy sender since p2p
117
                // connections don't have a bus to assign unique
118
                // names
119
                UniqueName::try_from(":p2p.test").unwrap()
120
            }
121
            #[cfg(not(any(test, feature = "test-util")))]
122
            {
123
                return Err(custom_service_error("Failed to get sender from header."));
124
            }
125
        };
126

            
127
170
        let peer_info = async {
128
63
            let proxy = zbus::fdo::DBusProxy::new(&self.connection()).await.ok()?;
129
142
            let pid = proxy
130
32
                .get_connection_unix_process_id(sender.as_ref().into())
131
124
                .await
132
                .ok()?;
133
            let cmdline = tokio::fs::read(format!("/proc/{pid}/cmdline")).await.ok()?;
134
            let name = cmdline.split(|&b| b == 0).next()?;
135
            let name = std::path::Path::new(std::str::from_utf8(name).ok()?)
136
                .file_name()?
137
                .to_str()?
138
                .to_owned();
139
            let session_type = SessionType::detect(pid).await;
140
            Some(PeerInfo::new(pid, name, session_type))
141
        }
142
125
        .await;
143
        let session = Session::new(
144
59
            aes_key.map(Arc::new),
145
61
            self.clone(),
146
31
            sender.clone(),
147
27
            peer_info,
148
        )
149
96
        .await;
150
61
        let path = OwnedObjectPath::from(session.path().clone());
151

            
152
57
        match session.peer_info() {
153
            Some(info) => {
154
                tracing::info!("Client {} ({}) connected, session: {}", sender, info, path)
155
            }
156
49
            None => tracing::info!("Client {} connected, session: {}", sender, path),
157
        }
158

            
159
156
        self.sessions
160
            .lock()
161
86
            .await
162
53
            .insert(path.clone(), session.clone());
163

            
164
38
        object_server.at(&path, session).await?;
165

            
166
24
        let service_key = public_key
167
30
            .map(OwnedValue::from)
168
83
            .unwrap_or_else(|| Value::new::<Vec<u8>>(vec![]).try_into_owned().unwrap());
169

            
170
28
        Ok((service_key, path))
171
    }
172

            
173
    #[zbus(out_args("collection", "prompt"))]
174
11
    pub async fn create_collection(
175
        &self,
176
        properties: Properties,
177
        alias: &str,
178
    ) -> Result<(OwnedObjectPath, ObjectPath<'_>), ServiceError> {
179
21
        let label = properties.label().to_owned();
180
21
        let alias = alias.to_owned();
181

            
182
        // Create a prompt to get the password for the new collection
183
        let prompt = Prompt::new(
184
21
            self.clone(),
185
            PromptRole::CreateCollection,
186
10
            label.clone(),
187
11
            None,
188
        )
189
34
        .await;
190
18
        let prompt_path = OwnedObjectPath::from(prompt.path().clone());
191

            
192
        // Store the collection metadata for later creation
193
40
        self.pending_collections
194
            .lock()
195
33
            .await
196
10
            .insert(prompt_path.clone(), (label, alias));
197

            
198
        // Create the collection creation action
199
21
        let service = self.clone();
200
11
        let creation_prompt_path = prompt_path.clone();
201
48
        let action = PromptAction::new(move |secret: Option<Secret>| async move {
202
42
            let collection_path = service
203
19
                .complete_collection_creation(&creation_prompt_path, secret)
204
43
                .await?;
205

            
206
20
            Ok(Value::new(collection_path).try_into_owned().unwrap())
207
        });
208

            
209
15
        prompt.set_action(action).await;
210

            
211
        // Register the prompt
212
55
        self.prompts
213
            .lock()
214
34
            .await
215
20
            .insert(prompt_path.clone(), prompt.clone());
216

            
217
15
        self.object_server().at(&prompt_path, prompt).await?;
218

            
219
18
        tracing::debug!("CreateCollection prompt created at `{}`", prompt_path);
220

            
221
        // Return empty collection path and the prompt path
222
23
        Ok((OwnedObjectPath::default(), prompt_path.into()))
223
    }
224

            
225
    #[zbus(out_args("unlocked", "locked"))]
226
4
    pub async fn search_items(
227
        &self,
228
        attributes: HashMap<String, String>,
229
    ) -> Result<(Vec<OwnedObjectPath>, Vec<OwnedObjectPath>), ServiceError> {
230
4
        let mut unlocked = Vec::new();
231
4
        let mut locked = Vec::new();
232
12
        let collections = self.collections.lock().await;
233

            
234
12
        for collection in collections.values() {
235
20
            let items = collection.search_inner_items(&attributes).await?;
236
8
            for item in items {
237
24
                if item.is_locked().await {
238
8
                    locked.push(item.path().clone().into());
239
                } else {
240
11
                    unlocked.push(item.path().clone().into());
241
                }
242
            }
243
        }
244

            
245
8
        if unlocked.is_empty() && locked.is_empty() {
246
8
            tracing::debug!(
247
                "Items with attributes {:?} does not exist in any collection.",
248
                attributes
249
            );
250
        } else {
251
8
            tracing::debug!("Items with attributes {:?} found.", attributes);
252
        }
253

            
254
4
        Ok((unlocked, locked))
255
    }
256

            
257
    #[zbus(out_args("unlocked", "prompt"))]
258
10
    pub async fn unlock(
259
        &self,
260
        objects: Vec<OwnedObjectPath>,
261
    ) -> Result<(Vec<OwnedObjectPath>, OwnedObjectPath), ServiceError> {
262
23
        tracing::info!("Unlock requested for {} objects.", objects.len());
263
24
        let (unlocked, not_unlocked) = self.set_locked(false, &objects).await?;
264
20
        if !not_unlocked.is_empty() {
265
            // Extract the label and collection before creating the prompt
266
20
            let label = self.extract_label_from_objects(&not_unlocked).await;
267
20
            let collection = self.extract_collection_from_objects(&not_unlocked).await;
268

            
269
20
            let prompt = Prompt::new(self.clone(), PromptRole::Unlock, label, collection).await;
270
16
            let path = OwnedObjectPath::from(prompt.path().clone());
271

            
272
            // Create the unlock action
273
8
            let service = self.clone();
274
59
            let action = PromptAction::new(move |secret: Option<Secret>| async move {
275
                // The prompter will handle secret validation
276
                // Here we just perform the unlock operation
277

            
278
                // First, check for pending migrations (without holding
279
                // collections lock)
280
32
                for object in &not_unlocked {
281
                    let collection = {
282
29
                        let collections = service.collections.lock().await;
283
18
                        collections.get(object).cloned()
284
                    };
285

            
286
9
                    if let Some(collection) = collection {
287
                        // Check if this collection has a pending migration by
288
                        // name
289
                        let migration_opt = {
290
30
                            let pending = service.pending_migrations.lock().await;
291
17
                            pending.get(collection.name()).cloned()
292
                        };
293

            
294
8
                        if let Some(migration) = migration_opt {
295
9
                            let migration_name = migration.name();
296
5
                            tracing::debug!(
297
                                "Attempting migration for '{}' during unlock",
298
                                migration_name
299
                            );
300

            
301
                            // Attempt migration with the provided secret (no
302
                            // locks held)
303
24
                            match migration.migrate(&service.data_dir, secret.as_ref()).await {
304
4
                                Ok(unlocked_keyring) => {
305
8
                                    tracing::info!(
306
                                        "Successfully migrated '{}' during unlock",
307
                                        migration_name
308
                                    );
309

            
310
                                    // Replace the keyring in the collection
311
12
                                    let mut keyring_guard = collection.keyring.write().await;
312
4
                                    *keyring_guard = Some(Keyring::Unlocked(unlocked_keyring));
313
4
                                    drop(keyring_guard);
314

            
315
                                    // Dispatch items from the migrated keyring
316
18
                                    if let Err(e) = collection.dispatch_items().await {
317
                                        tracing::error!(
318
                                            "Failed to dispatch items after migration: {}",
319
                                            e
320
                                        );
321
                                    }
322

            
323
                                    // Remove from pending migrations
324
23
                                    service
325
                                        .pending_migrations
326
6
                                        .lock()
327
28
                                        .await
328
6
                                        .remove(migration_name);
329
                                }
330
                                Err(e) => {
331
                                    tracing::warn!(
332
                                        "Failed to migrate '{}' during unlock: {}",
333
                                        migration_name,
334
                                        e
335
                                    );
336
                                    let _ = collection.set_locked(false, secret.clone()).await;
337
                                }
338
                            }
339
                        } else {
340
                            // Normal unlock
341
30
                            let _ = collection.set_locked(false, secret.clone()).await;
342
                        }
343
                    } else {
344
                        // Try to find as item within collections
345
16
                        let collections = service.collections.lock().await;
346
4
                        let mut found_collection = None;
347
8
                        for collection in collections.values() {
348
16
                            if let Some(item) = collection.item_from_path(object).await {
349
4
                                found_collection = Some((
350
4
                                    collection.clone(),
351
4
                                    item.clone(),
352
16
                                    collection.is_locked().await,
353
                                ));
354
                                break;
355
                            }
356
                        }
357
4
                        drop(collections);
358

            
359
4
                        if let Some((collection, item, is_locked)) = found_collection {
360
4
                            if is_locked {
361
16
                                let _ = collection.set_locked(false, secret.clone()).await;
362
                            } else {
363
                                let keyring = collection.keyring.read().await;
364
                                match keyring.as_ref() {
365
                                    Some(k) if !k.is_locked() => {
366
                                        let _ = item.set_locked(false, k.as_unlocked()).await;
367
                                    }
368
                                    _ => {
369
                                        drop(keyring);
370
                                        let _ = collection.set_locked(false, secret.clone()).await;
371
                                    }
372
                                }
373
                            }
374
                        }
375
                    }
376
                }
377
9
                Ok(Value::new(not_unlocked).try_into_owned().unwrap())
378
            });
379

            
380
12
            prompt.set_action(action).await;
381

            
382
40
            self.prompts
383
                .lock()
384
28
                .await
385
16
                .insert(path.clone(), prompt.clone());
386

            
387
12
            self.object_server().at(&path, prompt).await?;
388
8
            return Ok((unlocked, path));
389
        }
390

            
391
8
        Ok((unlocked, OwnedObjectPath::default()))
392
    }
393

            
394
    #[zbus(out_args("locked", "Prompt"))]
395
10
    pub async fn lock(
396
        &self,
397
        objects: Vec<OwnedObjectPath>,
398
    ) -> Result<(Vec<OwnedObjectPath>, OwnedObjectPath), ServiceError> {
399
23
        tracing::info!("Lock requested for {} objects.", objects.len());
400
27
        let (locked, not_locked) = self.set_locked(true, &objects).await?;
401
        // Locking never requires prompts, so not_locked should always be empty
402
        debug_assert!(
403
            not_locked.is_empty(),
404
            "Lock operation should never require prompts"
405
        );
406
10
        Ok((locked, OwnedObjectPath::default()))
407
    }
408

            
409
    #[zbus(out_args("secrets"))]
410
5
    pub async fn get_secrets(
411
        &self,
412
        items: Vec<OwnedObjectPath>,
413
        session: OwnedObjectPath,
414
    ) -> Result<HashMap<OwnedObjectPath, DBusSecretInner>, ServiceError> {
415
        tracing::debug!(
416
            "GetSecrets called for {} items with session {}.",
417
            items.len(),
418
            session
419
        );
420
5
        let mut secrets = HashMap::new();
421
15
        let collections = self.collections.lock().await;
422

            
423
20
        'outer: for collection in collections.values() {
424
20
            for item in &items {
425
16
                if let Some(item) = collection.item_from_path(item).await {
426
16
                    match item.get_secret(session.clone()).await {
427
4
                        Ok((secret,)) => {
428
8
                            secrets.insert(item.path().clone().into(), secret);
429
                            // To avoid iterating through all the remaining
430
                            // collections, if the
431
                            // items secrets are already retrieved.
432
4
                            if secrets.len() == items.len() {
433
                                break 'outer;
434
                            }
435
                        }
436
                        // Avoid erroring out if an item is locked.
437
                        Err(ServiceError::IsLocked(_)) => {
438
                            continue;
439
                        }
440
4
                        Err(err) => {
441
4
                            return Err(err);
442
                        }
443
                    };
444
                }
445
            }
446
        }
447

            
448
2
        tracing::debug!(
449
            "GetSecrets returned {} of {} requested secrets.",
450
            secrets.len(),
451
            items.len()
452
        );
453
4
        Ok(secrets)
454
    }
455

            
456
    #[zbus(out_args("collection"))]
457
92
    pub async fn read_alias(&self, name: &str) -> Result<OwnedObjectPath, ServiceError> {
458
        // Map "login" alias to "default" for compatibility with gnome-keyring
459
82
        let alias_to_find = if name == Self::LOGIN_ALIAS {
460
            oo7::dbus::Service::DEFAULT_COLLECTION
461
        } else {
462
26
            name
463
        };
464

            
465
30
        let collections = self.collections.lock().await;
466

            
467
103
        for (path, collection) in collections.iter() {
468
80
            if collection.alias().await == alias_to_find {
469
42
                tracing::debug!("Collection: {} found for alias: {}.", path, name);
470
47
                return Ok(path.to_owned());
471
            }
472
        }
473

            
474
8
        tracing::info!("Collection with alias {} does not exist.", name);
475

            
476
12
        Ok(OwnedObjectPath::default())
477
    }
478

            
479
4
    pub async fn set_alias(
480
        &self,
481
        name: &str,
482
        collection: OwnedObjectPath,
483
    ) -> Result<(), ServiceError> {
484
12
        let collections = self.collections.lock().await;
485

            
486
12
        for (path, other_collection) in collections.iter() {
487
8
            if *path == collection {
488
8
                other_collection.set_alias(name).await;
489

            
490
4
                tracing::info!("Collection: {} alias updated to {}.", collection, name);
491
4
                return Ok(());
492
            }
493
        }
494

            
495
4
        tracing::info!("Collection: {} does not exist.", collection);
496

            
497
8
        Err(ServiceError::NoSuchObject(format!(
498
            "The collection: {collection} does not exist.",
499
        )))
500
    }
501

            
502
    #[zbus(property, name = "Collections")]
503
112
    pub async fn collections(&self) -> Vec<OwnedObjectPath> {
504
74
        self.collections.lock().await.keys().cloned().collect()
505
    }
506

            
507
    #[zbus(signal, name = "CollectionCreated")]
508
    pub async fn collection_created(
509
10
        signal_emitter: &SignalEmitter<'_>,
510
10
        collection: &ObjectPath<'_>,
511
    ) -> zbus::Result<()>;
512

            
513
    #[zbus(signal, name = "CollectionDeleted")]
514
    pub async fn collection_deleted(
515
9
        signal_emitter: &SignalEmitter<'_>,
516
8
        collection: &ObjectPath<'_>,
517
    ) -> zbus::Result<()>;
518

            
519
    #[zbus(signal, name = "CollectionChanged")]
520
    pub async fn collection_changed(
521
10
        signal_emitter: &SignalEmitter<'_>,
522
10
        collection: &ObjectPath<'_>,
523
    ) -> zbus::Result<()>;
524
}
525

            
526
impl Service {
527
    const LOGIN_ALIAS: &str = "login";
528

            
529
    /// Set the prompter type override
530
    #[allow(unused)]
531
118
    pub(crate) async fn set_prompter_type(&self, prompter_type: PrompterType) {
532
72
        *self.prompter_type_override.lock().await = Some(prompter_type);
533
    }
534

            
535
    /// Get the prompter type to use based on the caller's session type
536
47
    pub(crate) async fn prompter_type(&self, peer_info: Option<&PeerInfo>) -> PrompterType {
537
41
        if let Some(override_type) = self.prompter_type_override.lock().await.as_ref() {
538
12
            return *override_type;
539
        }
540

            
541
        let is_graphical = peer_info
542
            .map(|p| p.session_type().is_graphical())
543
            .unwrap_or_default();
544

            
545
        if is_graphical {
546
            #[cfg(any(feature = "plasma_native_crypto", feature = "plasma_openssl_crypto"))]
547
            {
548
                if in_plasma_environment(&self.connection()).await {
549
                    return PrompterType::Plasma;
550
                }
551
            }
552

            
553
            return PrompterType::GNOME;
554
        }
555

            
556
        PrompterType::Cli
557
    }
558

            
559
32
    pub(crate) fn new(
560
        data_dir: std::path::PathBuf,
561
        pam_socket: Option<std::path::PathBuf>,
562
    ) -> Self {
563
        Self {
564
63
            collections: Arc::new(Mutex::new(HashMap::new())),
565
65
            connection: Default::default(),
566
67
            sessions: Arc::new(Mutex::new(HashMap::new())),
567
68
            session_index: Arc::new(AtomicU32::new(0)),
568
70
            prompts: Arc::new(Mutex::new(HashMap::new())),
569
69
            prompt_index: Arc::new(AtomicU32::new(0)),
570
68
            pending_collections: Arc::new(Mutex::new(HashMap::new())),
571
67
            pending_migrations: Arc::new(Mutex::new(HashMap::new())),
572
            data_dir,
573
            pam_socket,
574
67
            prompter_type_override: Arc::new(Mutex::new(None)),
575
        }
576
    }
577

            
578
    pub async fn run(
579
        secret: Option<Secret>,
580
        request_replacement: bool,
581
    ) -> Result<zbus::Connection, Error> {
582
        // Compute data directory from environment variables
583
        let data_dir = std::env::var_os("XDG_DATA_HOME")
584
            .filter(|h| !h.is_empty())
585
            .map(std::path::PathBuf::from)
586
            .filter(|p| p.is_absolute())
587
            .or_else(|| {
588
                std::env::var_os("HOME")
589
                    .filter(|h| !h.is_empty())
590
                    .map(std::path::PathBuf::from)
591
                    .map(|p| p.join(".local/share"))
592
            })
593
            .ok_or_else(|| {
594
                Error::IO(std::io::Error::new(
595
                    std::io::ErrorKind::NotFound,
596
                    "No data directory found (XDG_DATA_HOME or HOME)",
597
                ))
598
            })?;
599

            
600
        // Compute PAM socket path from environment variable
601
        let pam_socket = std::env::var_os("OO7_PAM_SOCKET").map(std::path::PathBuf::from);
602

            
603
        let service = Self::new(data_dir, pam_socket);
604

            
605
        // Start PAM listener early so it can buffer secrets arriving before
606
        // D-Bus is ready (e.g. during PAM-initiated login startup).
607
        tracing::info!("Starting PAM listener");
608
        let pam_listener = crate::pam_listener::PamListener::new(service.clone());
609
        let pam_listener_replay = pam_listener.clone();
610
        tokio::spawn(async move {
611
            if let Err(e) = pam_listener.start().await {
612
                tracing::error!("PAM listener error: {}", e);
613
            }
614
        });
615

            
616
        // The name is requested only once the service is initialized
617
        let connection = zbus::connection::Builder::session()?
618
            .serve_at(
619
                oo7::dbus::api::Service::PATH.as_deref().unwrap(),
620
                service.clone(),
621
            )?
622
            .build()
623
            .await?;
624

            
625
        #[cfg(any(feature = "gnome_native_crypto", feature = "gnome_openssl_crypto"))]
626
        connection
627
            .object_server()
628
            .at(
629
                oo7::dbus::api::Service::PATH.as_deref().unwrap(),
630
                InternalInterface::new(service.clone()),
631
            )
632
            .await?;
633

            
634
        // Discover existing keyrings
635
        let discovered_keyrings = service.discover_keyrings(secret.clone()).await?;
636

            
637
        let connection_clone = connection.clone();
638
        service
639
            .initialize(connection, discovered_keyrings, secret, true)
640
            .await?;
641

            
642
        let mut flags = RequestNameFlags::AllowReplacement | RequestNameFlags::DoNotQueue;
643
        if request_replacement {
644
            flags |= RequestNameFlags::ReplaceExisting;
645
        }
646
        match connection_clone
647
            .request_name_with_flags(
648
                oo7::dbus::api::Service::DESTINATION.as_deref().unwrap(),
649
                flags,
650
            )
651
            .await?
652
        {
653
            RequestNameReply::PrimaryOwner | RequestNameReply::AlreadyOwner => {}
654
            RequestNameReply::Exists | RequestNameReply::InQueue => {
655
                return Err(Error::Zbus(zbus::Error::NameTaken));
656
            }
657
        }
658

            
659
        // Replay any secrets that the PAM listener buffered during startup
660
        pam_listener_replay.replay_buffered_secrets().await;
661

            
662
        Ok(connection_clone)
663
    }
664

            
665
    #[cfg(any(test, feature = "test-util"))]
666
    #[allow(dead_code)]
667
32
    pub async fn run_with_connection(
668
        connection: zbus::Connection,
669
        data_dir: std::path::PathBuf,
670
        pam_socket: Option<std::path::PathBuf>,
671
        secret: Option<Secret>,
672
    ) -> Result<Self, Error> {
673
30
        let service = Self::new(data_dir, pam_socket);
674

            
675
        // Serve the service at the standard path
676
139
        connection
677
            .object_server()
678
            .at(
679
33
                oo7::dbus::api::Service::PATH.as_deref().unwrap(),
680
34
                service.clone(),
681
            )
682
102
            .await?;
683

            
684
        #[cfg(any(feature = "gnome_native_crypto", feature = "gnome_openssl_crypto"))]
685
123
        connection
686
            .object_server()
687
            .at(
688
35
                oo7::dbus::api::Service::PATH.as_deref().unwrap(),
689
28
                InternalInterface::new(service.clone()),
690
            )
691
110
            .await?;
692

            
693
63
        let default_keyring = if let Some(secret) = secret.clone() {
694
90
            vec![(
695
29
                "default".to_owned(),
696
35
                "Login".to_owned(),
697
29
                oo7::dbus::Service::DEFAULT_COLLECTION.to_owned(),
698
97
                Keyring::Unlocked(UnlockedKeyring::temporary(secret).await?),
699
            )]
700
        } else {
701
12
            vec![]
702
        };
703

            
704
128
        service
705
28
            .initialize(connection, default_keyring, secret, false)
706
120
            .await?;
707
27
        Ok(service)
708
    }
709

            
710
    /// Generate a unique label and alias by checking registered
711
    /// collections and appending a counter if needed. Returns a tuple of
712
    /// (label, alias).
713
30
    fn make_unique_label_and_alias(
714
        collections: &HashMap<OwnedObjectPath, Collection>,
715
        label: &str,
716
        alias: &str,
717
    ) -> (String, String) {
718
        // Sanitize the label to create the path (for checking uniqueness)
719
33
        let base_path = crate::collection::collection_path(label)
720
            .expect("Sanitized label should always produce valid object path");
721
64
        if !collections.contains_key(&base_path) {
722
68
            return (label.to_owned(), alias.to_owned());
723
        }
724

            
725
        // Append counter until we find a unique one
726
4
        let mut counter = 2;
727
4
        loop {
728
8
            let path = crate::collection::collection_path(&format!("{label}{counter}"))
729
                .expect("Sanitized label should always produce valid object path");
730
4
            let new_label = format!("{}{}", label, counter);
731
8
            let new_alias = format!("{}{}", alias, counter);
732

            
733
8
            if !collections.contains_key(&path) {
734
4
                return (new_label, new_alias);
735
            }
736
            counter += 1;
737
        }
738
    }
739

            
740
    /// Discover existing keyrings in the data directory
741
    /// Returns a vector of (name, label, alias, keyring) tuples
742
4
    pub(crate) async fn discover_keyrings(
743
        &self,
744
        secret: Option<Secret>,
745
    ) -> Result<Vec<(String, String, String, Keyring)>, Error> {
746
4
        let mut discovered = Vec::new();
747

            
748
8
        let keyrings_dir = self.data_dir.join("keyrings");
749

            
750
        // Scan for v1 keyrings first
751
8
        let v1_dir = keyrings_dir.join("v1");
752
8
        if v1_dir.exists() {
753
4
            tracing::debug!("Scanning for v1 keyrings in {}", v1_dir.display());
754
20
            if let Ok(mut entries) = tokio::fs::read_dir(&v1_dir).await {
755
20
                while let Ok(Some(entry)) = entries.next_entry().await {
756
4
                    let path = entry.path();
757

            
758
                    // Skip directories and non-.keyring files
759
8
                    if path.is_dir() || path.extension() != Some(std::ffi::OsStr::new("keyring")) {
760
                        continue;
761
                    }
762

            
763
12
                    if let Some(name) = path.file_stem().and_then(|s| s.to_str()) {
764
8
                        tracing::debug!("Found v1 keyring: {name}");
765

            
766
                        // Try to load the keyring
767
24
                        match self.load_keyring(&path, name, secret.as_ref()).await {
768
4
                            Ok((name, label, alias, keyring)) => {
769
4
                                discovered.push((name, label, alias, keyring))
770
                            }
771
                            Err(e) => tracing::warn!("Failed to load keyring {:?}: {}", path, e),
772
                        }
773
                    }
774
                }
775
            }
776
        }
777

            
778
        // Scan for v0 keyrings
779
8
        if keyrings_dir.exists() {
780
4
            tracing::debug!("Scanning for v0 keyrings in {}", keyrings_dir.display());
781
20
            if let Ok(mut entries) = tokio::fs::read_dir(&keyrings_dir).await {
782
20
                while let Ok(Some(entry)) = entries.next_entry().await {
783
4
                    let path = entry.path();
784

            
785
                    // Skip directories and non-.keyring files
786
8
                    if path.is_dir() || path.extension() != Some(std::ffi::OsStr::new("keyring")) {
787
                        continue;
788
                    }
789

            
790
4
                    if migration::stamp_path(&path).exists() {
791
                        continue;
792
                    }
793

            
794
12
                    if let Some(name) = path.file_stem().and_then(|s| s.to_str()) {
795
8
                        tracing::debug!("Found v0 keyring: {name}");
796

            
797
                        // Try to load the keyring
798
24
                        match self.load_keyring(&path, name, secret.as_ref()).await {
799
4
                            Ok((name, label, alias, keyring)) => {
800
4
                                discovered.push((name, label, alias, keyring))
801
                            }
802
4
                            Err(e) => tracing::warn!("Failed to load keyring {:?}: {}", path, e),
803
                        }
804
                    }
805
                }
806
            }
807
        }
808

            
809
        // Discover KWallet keyrings for migration
810
        #[cfg(feature = "kwallet_migration")]
811
        self.discover_kwallet_keyrings(&self.data_dir, secret.as_ref(), &mut discovered)
812
            .await;
813

            
814
12
        let pending_count = self.pending_migrations.lock().await.len();
815

            
816
8
        if discovered.is_empty() && pending_count == 0 {
817
8
            tracing::info!("No keyrings discovered in data directory");
818
        } else {
819
            tracing::info!(
820
                "Discovered {} keyring(s), {pending_count} pending migration(s)",
821
                discovered.len(),
822
            );
823
        }
824

            
825
4
        Ok(discovered)
826
    }
827

            
828
    /// Discover KWallet keyrings for migration
829
    #[cfg(feature = "kwallet_migration")]
830
    async fn discover_kwallet_keyrings(
831
        &self,
832
        data_dir: &std::path::Path,
833
        secret: Option<&Secret>,
834
        discovered: &mut Vec<(String, String, String, Keyring)>,
835
    ) {
836
        let kwallet_dir = data_dir.join("kwalletd");
837

            
838
        if !kwallet_dir.exists() {
839
            tracing::debug!("No kwalletd directory found, skipping KWallet discovery");
840
            return;
841
        }
842

            
843
        tracing::debug!("Scanning for KWallet files in {}", kwallet_dir.display());
844

            
845
        let Ok(mut entries) = tokio::fs::read_dir(&kwallet_dir).await else {
846
            tracing::warn!("Failed to read kwalletd directory");
847
            return;
848
        };
849

            
850
        while let Ok(Some(entry)) = entries.next_entry().await {
851
            let path = entry.path();
852

            
853
            // Only process .kwl files
854
            if path.extension().is_none_or(|ext| ext != "kwl") {
855
                continue;
856
            }
857

            
858
            if migration::stamp_path(&path).exists() {
859
                continue;
860
            }
861

            
862
            let Some(name) = path.file_stem().and_then(|s| s.to_str()) else {
863
                continue;
864
            };
865

            
866
            tracing::debug!("Found KWallet file: {name}");
867

            
868
            // Use lowercased name as alias
869
            let alias = name.to_lowercase();
870

            
871
            let label = {
872
                let mut chars = name.chars();
873
                match chars.next() {
874
                    None => String::new(),
875
                    Some(first) => first.to_uppercase().collect::<String>() + chars.as_str(),
876
                }
877
            };
878

            
879
            let migration = PendingMigration::KWallet {
880
                name: name.to_owned(),
881
                path: path.clone(),
882
                label: label.clone(),
883
                alias: alias.clone(),
884
            };
885

            
886
            if let Some(secret) = secret {
887
                tracing::debug!("Attempting immediate migration of KWallet keyring '{name}'",);
888
                match migration.migrate(&self.data_dir, Some(secret)).await {
889
                    Ok(unlocked) => {
890
                        tracing::info!("Successfully migrated KWallet keyring '{name}' to oo7",);
891
                        discovered.push((
892
                            name.to_owned(),
893
                            label,
894
                            alias,
895
                            Keyring::Unlocked(unlocked),
896
                        ));
897
                        continue;
898
                    }
899
                    Err(e) => {
900
                        tracing::warn!(
901
                            "Failed to migrate KWallet keyring '{name}' at {}: {e}. Creating locked placeholder collection.",
902
                            migration.path().display()
903
                        );
904
                    }
905
                }
906
            }
907

            
908
            // Migration failed or no secret - create locked placeholder and
909
            // register for pending migration
910
            tracing::debug!(
911
                "Creating locked placeholder for KWallet keyring '{name}', will migrate on unlock",
912
            );
913

            
914
            match LockedKeyring::open_at(&self.data_dir, name).await {
915
                Ok(locked) => {
916
                    tracing::debug!(
917
                        "Created locked placeholder for '{name}', adding to pending migrations",
918
                    );
919
                    discovered.push((
920
                        name.to_owned(),
921
                        label.clone(),
922
                        alias.clone(),
923
                        Keyring::Locked(locked),
924
                    ));
925
                    self.pending_migrations
926
                        .lock()
927
                        .await
928
                        .insert(name.to_owned(), migration);
929
                }
930
                Err(e) => {
931
                    tracing::error!("Failed to create placeholder keyring for '{name}': {e}");
932
                }
933
            }
934
        }
935
    }
936

            
937
    /// Load a single keyring from a file path
938
    /// Returns (name, label, alias, keyring)
939
4
    async fn load_keyring(
940
        &self,
941
        path: &std::path::Path,
942
        name: &str,
943
        secret: Option<&Secret>,
944
    ) -> Result<(String, String, String, Keyring), Error> {
945
12
        let alias = if name.eq_ignore_ascii_case(Self::LOGIN_ALIAS) {
946
8
            oo7::dbus::Service::DEFAULT_COLLECTION.to_owned()
947
        } else {
948
8
            name.to_owned().to_lowercase()
949
        };
950

            
951
        // Use name as label (capitalized for consistency with Login)
952
        let label = {
953
8
            let mut chars = name.chars();
954
4
            match chars.next() {
955
                None => String::new(),
956
4
                Some(first) => first.to_uppercase().collect::<String>() + chars.as_str(),
957
            }
958
        };
959

            
960
        // Try to load the keyring
961
20
        let keyring = match LockedKeyring::load(path).await {
962
4
            Ok(locked_keyring) => {
963
                // Successfully loaded as v1 keyring
964
4
                if let Some(secret) = secret {
965
8
                    match locked_keyring.unlock(secret.clone()).await {
966
4
                        Ok(unlocked) => {
967
8
                            tracing::info!("Unlocked keyring '{}' from {:?}", name, path);
968
4
                            Keyring::Unlocked(unlocked)
969
                        }
970
4
                        Err(e) => {
971
8
                            tracing::warn!(
972
                                "Failed to unlock keyring '{}' with provided secret: {}. Keeping it locked.",
973
                                name,
974
                                e
975
                            );
976
                            // Reload as locked since unlock consumed it
977
16
                            Keyring::Locked(LockedKeyring::load(path).await?)
978
                        }
979
                    }
980
                } else {
981
8
                    tracing::debug!("No secret provided, keeping keyring '{}' locked", name);
982
4
                    Keyring::Locked(locked_keyring)
983
                }
984
            }
985
7
            Err(oo7::file::Error::VersionMismatch(Some(version)))
986
                if version.first() == Some(&0) =>
987
            // v0 is the legacy version
988
            {
989
                // This is a v0 keyring that needs migration
990
2
                tracing::info!(
991
                    "Found legacy v0 keyring '{name}' at {}, registering for migration",
992
                    path.display()
993
                );
994

            
995
                let migration = PendingMigration::V0 {
996
4
                    name: name.to_owned(),
997
4
                    path: path.to_path_buf(),
998
4
                    label: label.clone(),
999
4
                    alias: alias.clone(),
                };
8
                tracing::debug!("Attempting immediate migration of v0 keyring '{name}'",);
16
                match migration.migrate(&self.data_dir, secret).await {
4
                    Ok(unlocked) => {
8
                        tracing::info!("Successfully migrated v0 keyring '{name}' to v1",);
8
                        return Ok((name.to_owned(), label, alias, Keyring::Unlocked(unlocked)));
                    }
4
                    Err(e) => {
8
                        tracing::warn!(
                            "Failed to migrate v0 keyring '{name}': {e}. Creating locked placeholder collection.",
                        );
                    }
                }
                // Migration failed - create locked placeholder and register for
                // pending migration
4
                tracing::debug!(
                    "Creating locked placeholder for v0 keyring '{}', will migrate on unlock",
                    name
                );
16
                let locked = LockedKeyring::open(name).await?;
16
                self.pending_migrations
                    .lock()
16
                    .await
4
                    .insert(name.to_owned(), migration);
4
                Keyring::Locked(locked)
            }
            // Plain legacy keyrings, written by gnome-keyring for an empty
            // password, don't have the binary file header. Their migration
            // doesn't need a secret, so if it fails this isn't a keyring.
            Err(oo7::file::Error::FileHeaderMismatch(_)) => {
6
                tracing::info!(
                    "Found legacy plain keyring '{name}' at {}, migrating",
                    path.display()
                );
                let migration = PendingMigration::V0 {
4
                    name: name.to_owned(),
4
                    path: path.to_path_buf(),
4
                    label: label.clone(),
4
                    alias: alias.clone(),
                };
16
                Keyring::Unlocked(migration.migrate(&self.data_dir, secret).await?)
            }
            Err(e) => {
                return Err(e.into());
            }
        };
8
        Ok((name.to_owned(), label, alias, keyring))
    }
    /// Initialize the service with collections and start client disconnect
    /// handler
28
    pub(crate) async fn initialize(
        &self,
        connection: zbus::Connection,
        mut discovered_keyrings: Vec<(String, String, String, Keyring)>, /* (name, label, alias,
                                                                          * keyring) */
        secret: Option<Secret>,
        auto_create_default: bool,
    ) -> Result<(), Error> {
55
        *self.connection.write().unwrap() = Some(connection.clone());
30
        let object_server = connection.object_server();
63
        let mut collections = self.collections.lock().await;
        // Check if we have a default collection
115
        let has_default = discovered_keyrings.iter().any(|(_, _, alias, _)| {
32
            alias == oo7::dbus::Service::DEFAULT_COLLECTION || alias == Self::LOGIN_ALIAS
        });
28
        if !has_default && auto_create_default {
            tracing::info!("No default collection found, creating 'Login' keyring");
            let keyring = if let Some(secret) = secret {
                UnlockedKeyring::open_at(&self.data_dir, Self::LOGIN_ALIAS, Some(secret))
                    .await
                    .map(Keyring::Unlocked)
            } else {
                LockedKeyring::open_at(&self.data_dir, Self::LOGIN_ALIAS)
                    .await
                    .map(Keyring::Locked)
            };
            let keyring = keyring.inspect_err(|e| {
                tracing::error!("Failed to create default Login keyring: {}", e);
            })?;
            let is_locked = if keyring.is_locked() {
                "locked"
            } else {
                "unlocked"
            };
            discovered_keyrings.push((
                Self::LOGIN_ALIAS.to_owned(),
                "Login".to_owned(),
                oo7::dbus::Service::DEFAULT_COLLECTION.to_owned(),
                keyring,
            ));
            tracing::info!("Created default 'Login' collection ({})", is_locked);
        }
        // Build all collections under the lock, then register on D-Bus after
        // releasing it to avoid deadlocks with incoming calls.
32
        let mut built_collections: Vec<(Collection, String)> = Vec::new();
120
        for (name, label, alias, keyring) in discovered_keyrings {
68
            tracing::info!("Setting up collection '{name}' (alias: {alias}).");
49
            let (unique_label, unique_alias) =
                Self::make_unique_label_and_alias(&collections, &label, &alias);
77
            let collection =
                Collection::new(&name, &unique_label, &unique_alias, self.clone(), keyring).await;
60
            collections.insert(collection.path().to_owned().into(), collection.clone());
32
            built_collections.push((collection, unique_alias));
        }
        // Always create session collection (always temporary)
        let session_collection = Collection::new(
            "session",
            "session",
29
            oo7::dbus::Service::SESSION_COLLECTION,
66
            self.clone(),
104
            Keyring::Unlocked(UnlockedKeyring::temporary(Secret::random().unwrap()).await?),
        )
91
        .await;
58
        collections.insert(
58
            session_collection.path().to_owned().into(),
25
            session_collection.clone(),
        );
19
        drop(collections);
        // Now register on D-Bus and dispatch items without holding the lock
85
        for (collection, alias) in &built_collections {
120
            collection.dispatch_items().await?;
128
            object_server
31
                .at(collection.path(), collection.clone())
131
                .await?;
59
            if alias == oo7::dbus::Service::DEFAULT_COLLECTION {
159
                object_server
59
                    .at(DEFAULT_COLLECTION_ALIAS_PATH, collection.clone())
134
                    .await?;
            }
        }
30
        let session_path = session_collection.path().to_owned();
67
        object_server.at(&session_path, session_collection).await?;
        // Spawn client disconnect handler
29
        let service = self.clone();
91
        tokio::spawn(async move { service.on_client_disconnect().await });
        // Spawn stale session cleanup task
27
        let service = self.clone();
92
        tokio::spawn(async move { service.cleanup_stale_sessions().await });
27
        Ok(())
    }
111
    async fn on_client_disconnect(&self) -> zbus::Result<()> {
148
        let rule = zbus::MatchRule::builder()
24
            .msg_type(zbus::message::Type::Signal)
            .sender("org.freedesktop.DBus")?
            .interface("org.freedesktop.DBus")?
            .member("NameOwnerChanged")?
            .arg(2, "")?
            .build();
33
        let mut stream =
            zbus::MessageStream::for_match_rule(rule, &self.connection(), None).await?;
85
        while let Some(message) = stream.try_next().await? {
            let body = message.body();
            let Ok((_name, old_owner, new_owner)) =
                body.deserialize::<(String, Optional<UniqueName<'_>>, Optional<UniqueName<'_>>)>()
            else {
                continue;
            };
            debug_assert!(new_owner.is_none()); // We enforce that in the matching rule
            let old_owner = old_owner
                .as_ref()
                .expect("A disconnected client requires an old_owner");
            if let Some(session) = self.session_from_sender(old_owner).await {
                let client_name = match session.peer_info() {
                    Some(info) => format!("{old_owner} ({info})"),
                    None => old_owner.to_string(),
                };
                session.mark_stale().await;
                tracing::info!(
                    "Client {} disconnected. Session: {} marked for cleanup.",
                    client_name,
                    session.path()
                );
            }
        }
        Ok(())
    }
103
    async fn cleanup_stale_sessions(&self) {
52
        let mut interval = tokio::time::interval(std::time::Duration::from_secs(30));
        loop {
131
            interval.tick().await;
            let stale_paths: Vec<_> = {
63
                let sessions = self.sessions.lock().await;
28
                let mut paths = Vec::new();
76
                for (path, session) in sessions.iter() {
                    if session.is_stale().await {
                        paths.push(path.clone());
                    }
                }
28
                paths
            };
56
            for path in stale_paths {
                if let Some(session) = self.session(&path).await {
                    match session.close().await {
                        Ok(_) => tracing::info!("Stale session {} cleaned up.", path),
                        Err(err) => {
                            tracing::error!("Failed to clean up stale session {}: {}", path, err)
                        }
                    }
                }
            }
        }
    }
10
    pub async fn set_locked(
        &self,
        locked: bool,
        objects: &[OwnedObjectPath],
    ) -> Result<(Vec<OwnedObjectPath>, Vec<OwnedObjectPath>), ServiceError> {
10
        let mut without_prompt = Vec::new();
10
        let mut with_prompt = Vec::new();
24
        let collections = self.collections.lock().await;
40
        for object in objects {
50
            let resolved = Self::resolve_alias(&collections, object)
44
                .await
40
                .unwrap_or_else(|| object.clone());
10
            let mut found = false;
20
            for (path, collection) in collections.iter() {
24
                let collection_locked = collection.is_locked().await;
10
                if resolved == *path {
8
                    found = true;
8
                    if collection_locked == locked {
2
                        tracing::debug!(
                            "Collection: {} is already {}.",
                            resolved,
                            if locked { "locked" } else { "unlocked" }
                        );
12
                        without_prompt.push(resolved.clone());
8
                    } else if locked {
                        // Locking never requires a prompt
30
                        collection.set_locked(true, None).await?;
8
                        without_prompt.push(resolved.clone());
                    } else {
                        // Unlocking may require a prompt
16
                        with_prompt.push(resolved.clone());
                    }
                    break;
31
                } else if let Some(item) = collection.item_from_path(&resolved).await {
8
                    found = true;
                    // If collection is locked, can't perform any item
                    // lock/unlock operations
8
                    if collection_locked {
                        // Unlocking an item when collection is locked requires
                        // unlocking collection
6
                        if !locked {
8
                            with_prompt.push(resolved.clone());
                        } else {
                            // Can't lock an item when collection is locked
4
                            return Err(ServiceError::IsLocked(format!(
                                "Cannot lock item {} when collection is locked",
                                resolved
                            )));
                        }
28
                    } else if locked == item.is_locked().await {
5
                        tracing::debug!(
                            "Item: {} is already {}.",
                            resolved,
                            if locked { "locked" } else { "unlocked" }
                        );
12
                        without_prompt.push(resolved.clone());
                    } else {
24
                        let keyring = collection.keyring.read().await;
8
                        match keyring.as_ref() {
16
                            Some(k) if !k.is_locked() => {
28
                                item.set_locked(locked, k.as_unlocked()).await?;
8
                                without_prompt.push(resolved.clone());
                            }
                            _ => {
                                if locked {
                                    return Err(ServiceError::IsLocked(format!(
                                        "Cannot lock item {} when collection is locked",
                                        resolved
                                    )));
                                } else {
                                    with_prompt.push(resolved.clone());
                                }
                            }
                        }
                    }
                    break;
                }
            }
10
            if !found {
8
                tracing::warn!("Object: {} does not exist.", object);
            }
        }
10
        Ok((without_prompt, with_prompt))
    }
31
    pub fn connection(&self) -> zbus::Connection {
63
        self.connection
            .read()
            .unwrap()
            .clone()
            .expect("service is not initialized")
    }
32
    pub fn object_server(&self) -> zbus::ObjectServer {
31
        self.connection().object_server().clone()
    }
    /// Drop the service's reference to its connection.
    ///
    /// The connection's object server holds the service, which in turn holds
    /// the connection, so neither is ever freed otherwise.
    #[cfg(any(test, feature = "test-util"))]
    #[allow(dead_code)]
15
    pub(crate) fn release_connection(&self) {
15
        self.connection.write().unwrap().take();
    }
10
    async fn resolve_alias(
        collections: &HashMap<OwnedObjectPath, Collection>,
        path: &ObjectPath<'_>,
    ) -> Option<OwnedObjectPath> {
30
        let alias = path.strip_prefix("/org/freedesktop/secrets/aliases/")?;
        let alias_to_find = if alias == Self::LOGIN_ALIAS {
            oo7::dbus::Service::DEFAULT_COLLECTION
        } else {
            alias
        };
        for (real_path, collection) in collections.iter() {
            if collection.alias().await == alias_to_find {
                return Some(real_path.clone());
            }
        }
        None
    }
60
    pub async fn collection_from_path(&self, path: &ObjectPath<'_>) -> Option<Collection> {
34
        let collections = self.collections.lock().await;
30
        if let Some(collection) = collections.get(path).cloned() {
15
            return Some(collection);
        }
16
        let resolved = Self::resolve_alias(&collections, path).await?;
        collections.get(&resolved).cloned()
    }
31
    pub fn session_index(&self) -> u32 {
25
        self.session_index.fetch_add(1, Ordering::Relaxed)
    }
    pub(crate) async fn session_from_sender(&self, sender: &UniqueName<'_>) -> Option<Session> {
        let sessions = self.sessions.lock().await;
        sessions.values().find(|s| s.sender() == sender).cloned()
    }
    pub async fn peer_display_name(&self, sender: &UniqueName<'_>) -> String {
        match self.session_from_sender(sender).await {
            Some(session) => match session.peer_info() {
                Some(info) => format!("{sender} ({info})"),
                None => sender.to_string(),
            },
            None => sender.to_string(),
        }
    }
72
    pub async fn session(&self, path: &ObjectPath<'_>) -> Option<Session> {
44
        let session = self.sessions.lock().await.get(path).cloned()?;
21
        session.unmark_stale().await;
19
        Some(session)
    }
24
    pub async fn remove_session(&self, path: &ObjectPath<'_>) {
16
        self.sessions.lock().await.remove(path);
    }
35
    pub async fn remove_collection(&self, path: &ObjectPath<'_>) {
24
        self.collections.lock().await.remove(path);
13
        if let Ok(signal_emitter) =
            self.signal_emitter(oo7::dbus::api::Service::PATH.as_deref().unwrap())
        {
21
            let _ = self.collections_changed(&signal_emitter).await;
        }
    }
12
    pub fn prompt_index(&self) -> u32 {
12
        self.prompt_index.fetch_add(1, Ordering::Relaxed)
    }
45
    pub async fn prompt(&self, path: &ObjectPath<'_>) -> Option<Prompt> {
26
        self.prompts.lock().await.get(path).cloned()
    }
42
    pub async fn remove_prompt(&self, path: &ObjectPath<'_>) {
24
        self.prompts.lock().await.remove(path);
        // Also clean up pending collection if it exists
14
        self.pending_collections.lock().await.remove(path);
    }
16
    pub async fn register_prompt(&self, path: OwnedObjectPath, prompt: Prompt) {
12
        self.prompts.lock().await.insert(path, prompt);
    }
11
    pub async fn pending_collection(
        &self,
        prompt_path: &ObjectPath<'_>,
    ) -> Option<(String, String)> {
44
        self.pending_collections
            .lock()
28
            .await
8
            .get(prompt_path)
            .cloned()
    }
11
    pub async fn create_collection_with_secret(
        &self,
        label: &str,
        alias: &str,
        secret: Option<Secret>,
    ) -> Result<OwnedObjectPath, ServiceError> {
        // Create a persistent keyring with the provided secret
71
        let keyring = UnlockedKeyring::open_at(&self.data_dir, &label.to_lowercase(), secret)
28
            .await
16
            .map_err(|err| custom_service_error(&format!("Failed to create keyring: {err}")))?;
        // Write the keyring file to disk immediately
55
        keyring
            .write()
39
            .await
11
            .map_err(|err| custom_service_error(&format!("Failed to write keyring file: {err}")))?;
11
        let keyring = Keyring::Unlocked(keyring);
11
        let name = label.to_lowercase();
        // Create the collection with unique label and alias
11
        let (unique_label, unique_alias) = {
26
            let collections = self.collections.lock().await;
22
            Self::make_unique_label_and_alias(&collections, label, alias)
        };
35
        let collection =
            Collection::new(&name, &unique_label, &unique_alias, self.clone(), keyring).await;
19
        let collection_path: OwnedObjectPath = collection.path().to_owned().into();
        // Register with object server
52
        self.object_server()
20
            .at(collection.path(), collection.clone())
40
            .await?;
        // Add to collections
44
        self.collections
            .lock()
31
            .await
9
            .insert(collection_path.clone(), collection);
        // Emit CollectionCreated signal
11
        let service_path = oo7::dbus::api::Service::PATH.as_ref().unwrap();
9
        let signal_emitter = self.signal_emitter(service_path)?;
24
        Service::collection_created(&signal_emitter, &collection_path).await?;
        // Emit PropertiesChanged for Collections property to invalidate client
        // cache
14
        self.collections_changed(&signal_emitter).await?;
15
        tracing::info!(
            "Collection `{}` created with label '{}'",
            collection_path,
            label
        );
9
        Ok(collection_path)
    }
8
    pub async fn complete_collection_creation(
        &self,
        prompt_path: &ObjectPath<'_>,
        secret: Option<Secret>,
    ) -> Result<OwnedObjectPath, ServiceError> {
20
        let Some((label, alias)) = self.pending_collection(prompt_path).await else {
8
            return Err(ServiceError::NoSuchObject(format!(
                "No pending collection for prompt `{prompt_path}`"
            )));
        };
43
        let collection_path = self
19
            .create_collection_with_secret(&label, &alias, secret)
44
            .await?;
27
        self.pending_collections.lock().await.remove(prompt_path);
8
        Ok(collection_path)
    }
53
    pub fn signal_emitter<'a, P>(
        &self,
        path: P,
    ) -> Result<zbus::object_server::SignalEmitter<'a>, oo7::dbus::ServiceError>
    where
        P: TryInto<ObjectPath<'a>>,
        P::Error: Into<zbus::Error>,
    {
67
        let signal_emitter = zbus::object_server::SignalEmitter::new(&self.connection(), path)?;
48
        Ok(signal_emitter)
    }
    /// Extract the collection label from a list of object paths
    /// The objects can be either collections or items
32
    async fn extract_label_from_objects(&self, objects: &[OwnedObjectPath]) -> String {
16
        if objects.is_empty() {
            return String::new();
        }
        // Check if at least one of the objects is a Collection
32
        for object in objects {
28
            if let Some(collection) = self.collection_from_path(object).await {
20
                return collection.label().await;
            }
        }
        // Get the collection path from the first item
        // assumes all items are from the same collection
12
        if let Some(path_str) = objects.first().and_then(|p| p.as_str().rsplit_once('/')) {
4
            let collection_path = path_str.0;
12
            if let Ok(obj_path) = ObjectPath::try_from(collection_path)
12
                && let Some(collection) = self.collection_from_path(&obj_path).await
            {
12
                return collection.label().await;
            }
        }
        String::new()
    }
    /// Extract the collection from a list of object paths
    /// The objects can be either collections or items
8
    async fn extract_collection_from_objects(
        &self,
        objects: &[OwnedObjectPath],
    ) -> Option<Collection> {
16
        if objects.is_empty() {
            return None;
        }
        // Check if at least one of the objects is a Collection
32
        for object in objects {
28
            if let Some(collection) = self.collection_from_path(object).await {
8
                return Some(collection);
            }
        }
        // Get the collection path from the first item
        // (assumes all items are from the same collection)
12
        let path = objects
            .first()
            .unwrap()
            .as_str()
            .rsplit_once('/')
8
            .map(|(parent, _)| parent)?;
12
        self.collection_from_path(&ObjectPath::try_from(path).unwrap())
12
            .await
    }
    /// Attempt to migrate pending keyrings with the provided secret
    /// Returns a list of successfully migrated keyring names
28
    pub async fn migrate_pending_keyrings(&self, secret: &Secret) -> Vec<String> {
4
        let mut migrated = Vec::new();
12
        let mut pending = self.pending_migrations.lock().await;
4
        let mut to_remove = Vec::new();
16
        for (name, migration) in pending.iter() {
8
            tracing::debug!("Attempting to migrate pending keyring: {name}");
20
            match migration.migrate(&self.data_dir, Some(secret)).await {
4
                Ok(unlocked) => {
4
                    let label = migration.label();
4
                    let alias = migration.alias();
                    // Create a collection for this migrated keyring with unique
                    // label and alias
4
                    let (unique_label, unique_alias) = {
12
                        let collections = self.collections.lock().await;
8
                        Self::make_unique_label_and_alias(&collections, label, alias)
                    };
4
                    let keyring = Keyring::Unlocked(unlocked);
12
                    let collection =
                        Collection::new(name, &unique_label, &unique_alias, self.clone(), keyring)
24
                            .await;
4
                    let collection_path: OwnedObjectPath = collection.path().to_owned().into();
                    // Dispatch items
16
                    if let Err(e) = collection.dispatch_items().await {
                        tracing::error!(
                            "Failed to dispatch items for migrated keyring '{name}': {e}",
                        );
                        continue;
                    }
20
                    if let Err(e) = self
                        .object_server()
8
                        .at(collection.path(), collection.clone())
20
                        .await
                    {
                        tracing::error!(
                            "Failed to register migrated collection '{name}' with object server: {e}",
                        );
                        continue;
                    }
20
                    self.collections
                        .lock()
20
                        .await
8
                        .insert(collection_path.clone(), collection.clone());
4
                    if alias == oo7::dbus::Service::DEFAULT_COLLECTION
                        && let Err(e) = self
                            .object_server()
                            .at(DEFAULT_COLLECTION_ALIAS_PATH, collection)
                            .await
                    {
                        tracing::error!(
                            "Failed to register default alias for migrated collection '{name}': {e}",
                        );
                    }
8
                    if let Ok(signal_emitter) =
                        self.signal_emitter(oo7::dbus::api::Service::PATH.as_ref().unwrap())
                    {
12
                        let _ =
                            Service::collection_created(&signal_emitter, &collection_path).await;
16
                        let _ = self.collections_changed(&signal_emitter).await;
                    }
8
                    tracing::info!("Migrated keyring '{name}' added as collection",);
8
                    migrated.push(name.clone());
4
                    to_remove.push(name.clone());
                }
                Err(e) => {
                    tracing::debug!(
                        "Failed to migrate keyring '{name}' found at {} with provided secret: {e}",
                        migration.path().display()
                    );
                }
            }
        }
4
        for name in &to_remove {
8
            pending.remove(name);
        }
4
        migrated
    }
}
#[cfg(test)]
mod tests;
#[cfg(test)]
mod unencrypted_tests;